Elasticsearch 学习笔记

Elasticsearch 写入时先将每篇文章分词,再反向建立 “单个词汇→包含该词汇的所有文章” 的倒排索引,同时对词汇排序以支撑后续高效查询;搜索时则先借助内存中的 Term Index 前缀定位 + 二分查找快速找到目标词汇对应的文章 ID 列表,再根据 AND/OR 等搜索条件对多个词汇的文章 ID 列表取交集或并集,最后依据筛选后的文章 ID 提取完整内容并返回。

https://www.bilibili.com/video/BV1yb421J7oX

一、顶层设计:为什么需要 Elasticsearch?

1.1 ES 的定位

7-Blog/后端与微服务/assets/Elasticsearch_在系统中的定位-9b2df756

1.2 传统数据库面临的问题

7-Blog/后端与微服务/assets/传统关系型数据库的瓶颈-2e271bb1
  • 模糊查询性能瓶颈:MySQL 等关系型数据库使用 LIKE %keyword% 进行左模糊或全模糊查询时,无法利用索引(Index),会导致全表扫描,在百万级数据量下性能急剧下降。
  • 复杂维度筛选困难:当面临海量数据的多维度筛选、统计、标签聚合时,SQL 语句变得极度复杂且执行效率低下。
  • 文本相关性缺失:传统数据库只能做“精准匹配”或“简单的包含匹配”,无法计算相关性得分(Score),无法按“搜索结果匹配度”对结果进行排序。
  • 分词能力弱:无法处理同义词、纠错、中文分词(如将“由于”和“游泳”区分开)等复杂的自然语言处理需求。

1.3 Elasticsearch 解决的核心问题

7-Blog/后端与微服务/assets/Elasticsearch_核心价值-0a0cdb5f
  • 性能全文检索:基于倒排索引,实现亿级数据毫秒级响应。
  • 高可用与横向扩展:天生的分布式架构,通过增加节点即可线性提升存储容量和计算能力。
  • 复杂聚合分析:提供强大的聚合(Aggregations)框架,能够替代部分 OLAP 场景,实时统计数据分布。
  • 相关性排序:基于 TF-IDF / BM25 算法,提供符合人类直觉的搜索结果排名。

二、核心原理:如何解决问题?

2.1 倒排索引(Inverted Index)

这是 ES 快如闪电的根本原因。

  • 正排索引(Forward Index):文档 ID -> 文档内容(类似 MySQL 主键查询)。
  • 倒排索引:单词(Term)-> 包含该单词的文档 ID 列表。
    • ES 在写入数据时,会通过**分词器(Analyzer)**将文本拆解为单词,建立索引。
    • 查询时,直接根据单词找到文档 ID 列表,并通过位运算快速合并结果,无需扫描全表。
7-Blog/后端与微服务/assets/倒排索引原理详解-5d950c3c

2.2 Term Index(词项索引)

7-Blog/后端与微服务/assets/Term_Index_加速原理-fe766372
  • 面临的问题:Term Dictionary(词项字典)数据量极大(千万/亿级),无法全部放入内存,必须存在磁盘。每次查询都去磁盘翻字典,I/O 消耗巨大。
  • 解决方案Term Index 是一个基于词项前缀构建的精简目录树(类似 Trie 树或 FST - Finite State Transducers)。
    • 内存驻留:它体积非常小,可以完全加载到 内存(RAM) 中。
    • 加速原理:查询时,先在内存中查 Term Index,找到该 Term 在磁盘 Term Dictionary 中的大概位置(Offset 偏移量),然后再去磁盘读取具体内容。
    • 类比:Term Dictionary 是字典本身(在磁盘),Term Index 是字典的“拼音首字母索引页”(在内存)。

2.3 存储结构

7-Blog/后端与微服务/assets/ES_存储结构详解-c3bd152d

当倒排索引帮我们找到文档 ID 后,我们还需要获取内容或进行排序。

  • Stored Fields(行式存储)
    • 用途:用于存储原始文档内容_source 字段)。
    • 特点:行式存储,适合展示完整信息,但不适合聚合分析(读取很多冗余数据)。
  • Doc Values(列式存储)
    • 用途:专用于排序(Sorting)和聚合(Aggregations)
    • 原理:空间换时间。将散落在不同文档中的同一个字段值,集中在一起列式存储。
    • 优势:排序或计算平均值时,CPU 可以连续读取内存地址,极大提升效率,且对操作系统文件缓存(OS Cache)非常友好。

2.4 Segment 与 Lucene

7-Blog/后端与微服务/assets/Segment_与_Lucene_架构-6f82b8e4
  • Lucene:ES 的底层核心库。一个 ES 的分片(Shard)本质上就是一个完整的 Lucene 索引
  • Segment(段)
    • 最小单元:Lucene 内部由多个 Segment 组成。一个 Segment 包含了倒排索引、Term Index、Stored Fields 等所有结构,是具备完整搜索功能的最小单元。
    • 不可变性(Immutable):Segment 一旦生成,不可修改
    • 写入与合并:新增数据会生成新的 Segment;删除数据只是打上 .del 标记(逻辑删除)。ES 会在后台自动进行 Segment Merging(段合并),将多个小 Segment 合并为大 Segment,同时物理剔除被标记删除的数据。

三、分布式架构设计

7-Blog/后端与微服务/assets/ES_分布式架构设计-01752588
  • Cluster(集群):由多个节点组成。
  • Node(节点):单个服务实例。
  • Index(索引):逻辑上的数据集合(类比 Database)。
  • Shard(分片):数据被切分成多个分片存储在不同节点,实现并行计算和存储扩展。
  • Replica(副本):分片的备份,用于提高可用性(HA)和读取吞吐量。

3.1 高性能优化:分片(Shard)

7-Blog/后端与微服务/assets/高性能优化:分片机制-6bcd256e
  • 问题:单个 Index 数据量过大(如 1TB),单机硬盘存不下,且搜索时单线程扫描太慢。
  • 解决方案分片机制
    • 将一个 Index 逻辑拆分为多个 Shard(分片)。
    • 每个 Shard 是一个独立的 Lucene 实例。
    • 优势:读写压力被分散到多个 Shard 上并行处理,极大提升吞吐量。

3.2 高扩展性优化:多节点(Node)

7-Blog/后端与微服务/assets/高扩展性优化:多节点部署-fa229e0c
  • 原理横向扩展(Scale Out)
  • 机制:当数据量增长,分片增多,单机 CPU/内存吃紧时,可以向集群中加入新的机器(Node)。
  • 自动平衡:ES 会自动感知新节点,并将部分 Shard 迁移过去,实现负载均衡。

3.3 高可用优化:副本(Replica)

7-Blog/后端与微服务/assets/高可用优化:副本机制-27535f46
  • 角色区分
    • Primary Shard(主分片):负责处理写入请求。
    • Replica Shard(副本分片):主分片的完整备份。
  • 机制
    • 读写分离:副本分片可以分担搜索(读)请求,提升查询并发量。
    • 故障转移(Failover):如果持有主分片的 Node 挂掉,集群会迅速选举一个副本分片升级为新的主分片,保证服务不中断。

3.4 节点角色分化

7-Blog/后端与微服务/assets/Node_角色分化-523fd620

在大型集群中,让每个节点“各司其职”效率更高。

  • Master Node(主节点):集群的大脑。负责索引创建/删除、维护集群状态(Cluster State)、管理节点加入/退出。
  • Data Node(数据节点):苦力。负责存储数据(Shard),执行耗费资源的 CRUD 和聚合操作。
  • Coordinate Node(协调节点):前台/路由。接收客户端请求,分发给 Data Node,并汇总最终结果(Scatter-Gather)。

3.5 去中心化协调机制

7-Blog/后端与微服务/assets/去中心化协调机制-0c561c9f
  • Raft 算法衍生:ES 内部实现了一套基于 Raft 改进的共识算法(Zen Discovery)。
  • 作用
    • 保证集群中所有节点对“集群状态”的认知是一致的。
    • 实现 Master 节点的选举。
    • 故障检测:节点之间互相 Ping,感知是否有节点掉线。

四、核心流程详解

4.1 写入流程

ES 的写入流程设计权衡了数据安全性写入高吞吐

1. 路由与转发

  • 客户端向任意节点发送写入请求(该节点暂时成为协调节点)。
  • 协调节点使用路由算法确定数据所属的主分片(Primary Shard)位置:
    • shard = hash(routing) % number_of_primary_shards
    • 注:routing 默认是文档 _id
  • 协调节点将请求转发给持有该主分片的 Data Node。

2. 主分片写入 (Primary Operation)

  • 写入内存缓冲区 (Memory Buffer):数据先写入内存 buffer,此时数据不可被搜索
  • 写入 Translog (Transaction Log):同时追加写入 Translog 文件(顺序写磁盘),防止断电丢失数据。
  • Refresh (准实时关键步骤):默认每 1 秒,ES 将 buffer 中的数据生成一个新的 Segment 文件(此时建立倒排索引),并清空 buffer。一旦生成 Segment,数据即可被搜索。这就是 ES 被称为“准实时(Near Real-Time, NRT)”的原因。

3. 同步副本 (Replication)

  • 主分片写入成功后,并行将请求发送给所有的 副本分片 (Replica Shards)
  • 副本分片执行相同的写入逻辑。

4. 响应客户端

  • 当所有在 ISR (In-Sync Replicas) 列表中的副本都反馈写入成功后,主分片向协调节点报告成功。
  • 协调节点向客户端返回“写入完成”。

技术深挖:Flush 操作 Translog 不会无限增长。当 Translog 达到阈值或每隔 30 分钟,ES 会触发 Flush 操作:

  1. 强制执行 Refresh。
  2. 将所有内存中的 Segment 强制 fsync 刷入物理磁盘。
  3. 清空 Translog。 这保证了数据的持久化存储。
7-Blog/后端与微服务/assets/ES_写入流程详解-4fcf1ad9

4.2 搜索流程(Query Then Fetch)

7-Blog/后端与微服务/assets/ES_搜索流程详解-f639461c

搜索比写入复杂,因为数据分散在多个分片上,必须通过“两阶段”策略来整合结果,以避免网络带宽的巨大浪费。

阶段一:查询阶段 (Query Phase)

  • 请求分发:客户端向协调节点发送搜索请求。协调节点根据请求(是否有 routing 参数)将请求广播到所有相关分片(主分片或副本分片均可,负载均衡)。
  • 本地检索:每个分片在本地 Lucene 中执行搜索:
    1. 利用倒排索引筛选匹配文档。
    2. 利用 Doc Values 进行排序和打分。
    3. 关键点:分片仅返回文档 ID、相关性算分 (_score) 和排序值给协调节点,不返回文档的完整内容 (_source)。
  • 全局排序:协调节点收到所有分片返回的轻量级列表(例如每分片前 10 条),在内存中进行全局归并排序,选出最终的 Top N 文档 ID。

阶段二:获取阶段 (Fetch Phase)

  • 精确定位:协调节点知道了最终需要哪几个文档,以及它们位于哪个分片。
  • 抓取内容:协调节点向相关分片发送 Multi-Get 请求,只索取这 Top N 文档的完整内容 (_source / Stored Fields)。
  • 返回结果:分片返回文档详情,协调节点拼装最终 JSON,响应给客户端。

性能隐患:深度分页 (Deep Pagination) 如果查询 from=10000, size=10

  • 每个分片都必须查询出前 10010 条记录。
  • 假设有 5 个分片,协调节点需要接收 5 * 10010 = 50050 条记录的 ID,并在内存中排序,最后只取 10 条。
  • 后果:内存爆炸,CPU 飙升。
  • 对策:避免深分页,使用 Search AfterScroll API。

4.3 搜索流程总结图

7-Blog/后端与微服务/assets/搜索流程数据结构使用-f707625b

五、核心概念对比与 Type 演变

5.1 核心概念对比

┌─────────────────────────────────────────────────────────────────────────┐
│                      ES 与关系型数据库概念对比                              │
├───────────────────┬─────────────────────┬───────────────────────────────┤
│    关系型数据库     │    Elasticsearch    │           说明                │
├───────────────────┼─────────────────────┼───────────────────────────────┤
│    DatabaseCluster          │  数据库/集群                   │
│    TableIndex            │  表/索引                      │
│    RowDocument         │  行/文档                      │
│    ColumnField            │  列/字段                      │
│    SchemaMapping          │  表结构/映射                   │
│    IndexInverted Index   │  索引/倒排索引                 │
│    SQL            │    Query DSL        │  查询语言                      │
└───────────────────┴─────────────────────┴───────────────────────────────┘

5.2 Type 演变历史

7-Blog/后端与微服务/assets/Type_概念的演变-da19af3e
  • 5.x 及以前:允许一个 Index 下存在多个 Type(类比 Table),但本质上底层字段是扁平化混在一起的,导致数据稀疏(Sparse)问题,影响压缩效率和性能。
  • 6.x:强制规定一个 Index 只能有一个 Type,通常默认名为 doc
  • 7.x:Type 概念被彻底废弃(默认为 _doc),API 中 URL 的 type 参数变为可选。
  • 8.x:彻底移除 Type 概念。
    • 结论:现在设计索引时,严格遵循 “一个 Index 对应一类业务数据” 的原则。

为什么移除?

  • 映射爆炸(Mapping Explosion)
    • 在使用多类型时,如果不同类型之间有大量不同的字段,这会导致映射的数量急剧增加,进而引发映射爆炸问题。映射爆炸不仅会消耗大量的内存资源,还会降低 Elasticsearch 的性能,尤其是在处理大量数据时
  • 字段名冲突
    • 在同一个index的不同type中,如果有相同名称但映射类型不同的字段,会造成字段名冲突。这是因为 Elasticsearch 在内部是将这些字段扁平化处理的,而不同类型的相同名称字段可能会导致数据解析和查询时的混乱
    • 我们可以和关系型数据库来对比,在同一个数据库中,这些不同的表,可以有名称相同但类型不同的字段。而在 Elasticsearch 同一个index的不同type中,如果有不同document的字段名相同,但是类型不同,就会报错
  • 综上所述,Elasticsearch 从 7.x 版本开始废弃类型的主要目的是为了提升系统的性能、避免映射爆炸和字段冲突的问题,以及简化数据模型的设计和管理。这一改变反映了 Elasticsearch 对于提高性能、可维护性和用户体验的持续追求

六、核心功能

6.1 搜索能力矩阵

7-Blog/后端与微服务/assets/ES_搜索能力矩阵-cd6c4c21

6.2 聚合分析

7-Blog/后端与微服务/assets/聚合分析类型2-1609a946

七、应用场景

7-Blog/后端与微服务/assets/Elasticsearch_应用场景-39904ef5
  1. 日志与监控(ELK Stack):收集服务器日志、应用 Error 日志,快速定位故障(Logstash/Beats + ES + Kibana)。
  2. 站内搜索:电商商品搜索、论坛帖子搜索、企业知识库检索。
  3. 大屏可视化/BI:实时统计大盘数据(如双11大屏),利用聚合功能快速出报表。
  4. 地理位置服务(LBS):查询“附近的酒店”、“方圆5公里内的订单”(Geo-point/Geo-shape)。

八、如何使用 Elasticsearch

8.1 Spring Boot 集成

通常有两种主流方式:

  1. Spring Data Elasticsearch:封装程度极高,类似 JPA/MyBatis-Plus,通过 Repository 接口操作。
    • 优点:开发极快,代码简洁。
    • 缺点:灵活性稍差,对复杂 DSL 和版本兼容性控制不如原生客户端细致。
  2. RestHighLevelClient (官方推荐/传统):基于 HTTP 的原生客户端封装。
    • 优点:完全覆盖官方 API,灵活,可控性强。
    • 注意:ES 7.15+ 后官方推出了新的 Elasticsearch Java API Client,但 RestHighLevelClient 依然在存量系统中广泛使用。本文基于此方案进行封装。
<!-- Maven 依赖 -->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-elasticsearch</artifactId>
</dependency>
# application.yml 配置
spring:
  elasticsearch:
    uris: http://localhost:9200
    username: elastic
    password: password

此type的作用就是为了兼容6.x、7.x中的type概念,默认是关闭

8.2 原生 API 操作的复杂性

直接使用 RestHighLevelClient 会面临大量样板代码:

  • 构建 SearchSourceBuilderBoolQueryBuilder 极其繁琐。
  • 需要手动处理 IOException
  • 响应结果解析(Parse)需要从 JSON 层层剥离,非常痛苦。
  • 连接管理和配置分散。
// 原生 ES API 操作示例 - 复杂且繁琐
public SearchResponse searchPrograms(String keyword, Integer categoryId) {
    // 构建查询条件
    BoolQueryBuilder boolQuery = QueryBuilders.boolQuery();

    if (StringUtils.isNotBlank(keyword)) {
        boolQuery.must(QueryBuilders.matchQuery("title", keyword));
    }

    if (categoryId != null) {
        boolQuery.filter(QueryBuilders.termQuery("categoryId", categoryId));
    }

    // 构建排序
    FieldSortBuilder sortBuilder = SortBuilders.fieldSort("showTime")
            .order(SortOrder.ASC);

    // 构建搜索源
    SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
    sourceBuilder.query(boolQuery);
    sourceBuilder.sort(sortBuilder);
    sourceBuilder.from(0);
    sourceBuilder.size(10);

    // 构建搜索请求
    SearchRequest searchRequest = new SearchRequest("program-index");
    searchRequest.source(sourceBuilder);

    // 执行搜索
    return restHighLevelClient.search(searchRequest, RequestOptions.DEFAULT);
}

在平时开发中还是使用关系型数据库更加普遍,对于数据库、表、字段的概念更为熟悉,也更加习惯对表概念的操作。

而在操作Elasticsearch时,提供的api其实是很复杂的,各种操作的对象,如SearchSourceBuilder FieldSortBuilder BoolQueryBuilder 等等,操作上其实算不上简单,为了解决这个问题,在springboot操作Elasticsearch的基础上,进一步的封装,使用起来贴近于关系型数据库的方式,操作起来更加的容易上手

九、封装设计

9.1 封装目标

  • 统一配置:简化连接参数管理。
  • 屏蔽细节:隐藏繁琐的 Builder 构建过程。
  • 简化查询:通过 Map 或对象传递参数,自动构建 DSL。
  • 结果转换:自动将 ES 的 JSON 结果转为 Java Bean。
  • 健壮性:统一异常处理,防止 ES 波动导致服务崩溃。
7-Blog/后端与微服务/assets/封装设计目标-c8a9c448

9.2 封装架构设计

7-Blog/后端与微服务/assets/ES_封装框架架构-9d1b66fd
elasticsearch-spring-boot-starter/
├── src/main/java/com/example/elasticsearch/
│   ├── config/
│   │   ├── ElasticsearchProperties.java          # 配置属性
│   │   └── ElasticsearchAutoConfiguration.java   # 自动配置
│   ├── core/
│   │   ├── ElasticsearchService.java             # 核心服务类
│   │   ├── ElasticsearchIndexService.java        # 索引操作服务
│   │   └── ElasticsearchDocumentService.java     # 文档操作服务
│   ├── query/
│   │   ├── EsQueryBuilder.java                   # 查询构建器
│   │   ├── EsSearchRequest.java                  # 搜索请求封装
│   │   └── EsHighlightConfig.java                # 高亮配置
│   ├── result/
│   │   ├── PageResult.java                       # 分页结果
│   │   ├── SearchResult.java                     # 搜索结果
│   │   └── AggregationResult.java                # 聚合结果
│   ├── exception/
│   │   └── ElasticsearchException.java           # 自定义异常
│   └── annotation/
│       ├── EsDocument.java                       # 文档注解
│       └── EsField.java                          # 字段注解
└── src/main/resources/
    └── META-INF/spring.factories

9.3 配置类设计

package com.example.elasticsearch.config;

import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.validation.annotation.Validated;

import javax.validation.constraints.Min;
import javax.validation.constraints.NotEmpty;
import java.util.List;

/**
 * Elasticsearch 配置属性类
 * 支持集群配置、连接池、认证、重试等
 */
@Data
@Validated
@ConfigurationProperties(prefix = "elasticsearch")
public class ElasticsearchProperties {

    /**
     * 是否启用 ES
     */
    private Boolean enabled = true;

    /**
     * ES 节点地址列表(支持集群)
     */
    @NotEmpty(message = "ES 节点地址不能为空")
    private List<String> nodes = List.of("localhost:9200");

    /**
     * 用户名(可选)
     */
    private String username;

    /**
     * 密码(可选)
     */
    private String password;

    /**
     * 协议:http 或 https
     */
    private String scheme = "http";

    /**
     * 是否启用 Type(兼容 ES 6.x 版本)
     */
    private Boolean enableType = false;

    /**
     * 默认 Type 名称
     */
    private String defaultType = "_doc";

    /**
     * 连接超时时间(毫秒)
     */
    @Min(value = 1000, message = "连接超时时间不能小于1000ms")
    private Integer connectTimeout = 5000;

    /**
     * Socket 超时时间(毫秒)
     */
    @Min(value = 1000, message = "Socket超时时间不能小于1000ms")
    private Integer socketTimeout = 30000;

    /**
     * 请求超时时间(毫秒)
     */
    private Integer connectionRequestTimeout = 5000;

    /**
     * 最大连接数
     */
    @Min(value = 1, message = "最大连接数不能小于1")
    private Integer maxConnTotal = 100;

    /**
     * 每个路由的最大连接数
     */
    @Min(value = 1, message = "每个路由最大连接数不能小于1")
    private Integer maxConnPerRoute = 50;

    /**
     * 重试次数
     */
    @Min(value = 0, message = "重试次数不能为负数")
    private Integer retryTimes = 3;

    /**
     * 重试间隔(毫秒)
     */
    private Long retryInterval = 1000L;

    /**
     * 是否开启嗅探器
     */
    private Boolean enableSniffer = false;

    /**
     * 嗅探间隔时间(毫秒)
     */
    private Long snifferInterval = 60000L;

    /**
     * 批量操作每批大小
     */
    @Min(value = 100, message = "批量操作每批大小不能小于100")
    private Integer bulkBatchSize = 1000;

    /**
     * 批量操作刷新策略:immediate, wait_for, none
     */
    private String bulkRefreshPolicy = "none";

    /**
     * 是否打印 DSL 日志
     */
    private Boolean printDsl = false;

    /**
     * 慢查询阈值(毫秒),超过此值记录警告日志
     */
    private Long slowQueryThreshold = 3000L;
}

9.4 自动配置类

package com.example.elasticsearch.config;

import com.example.elasticsearch.core.ElasticsearchDocumentService;
import com.example.elasticsearch.core.ElasticsearchIndexService;
import com.example.elasticsearch.core.ElasticsearchService;
import lombok.extern.slf4j.Slf4j;
import org.apache.http.HttpHost;
import org.apache.http.auth.AuthScope;
import org.apache.http.auth.UsernamePasswordCredentials;
import org.apache.http.client.CredentialsProvider;
import org.apache.http.impl.client.BasicCredentialsProvider;
import org.elasticsearch.client.RestClient;
import org.elasticsearch.client.RestClientBuilder;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.client.sniff.Sniffer;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.util.StringUtils;

import javax.annotation.PreDestroy;
import java.util.List;
import java.util.stream.Collectors;

/**
 * Elasticsearch 自动配置类
 */
@Slf4j
@Configuration
@EnableConfigurationProperties(ElasticsearchProperties.class)
@ConditionalOnProperty(prefix = "elasticsearch", name = "enabled", havingValue = "true", matchIfMissing = true)
public class ElasticsearchAutoConfiguration {

    private RestHighLevelClient restHighLevelClient;
    private Sniffer sniffer;

    @Bean
    @ConditionalOnMissingBean
    public RestHighLevelClient restHighLevelClient(ElasticsearchProperties properties) {
        // 解析节点地址
        List<HttpHost> httpHosts = properties.getNodes().stream()
                .map(node -> {
                    String[] parts = node.split(":");
                    String host = parts[0];
                    int port = parts.length > 1 ? Integer.parseInt(parts[1]) : 9200;
                    return new HttpHost(host, port, properties.getScheme());
                })
                .collect(Collectors.toList());

        RestClientBuilder builder = RestClient.builder(
                httpHosts.toArray(new HttpHost[0])
        );

        // 设置请求配置
        builder.setRequestConfigCallback(requestConfigBuilder ->
                requestConfigBuilder
                        .setConnectTimeout(properties.getConnectTimeout())
                        .setSocketTimeout(properties.getSocketTimeout())
                        .setConnectionRequestTimeout(properties.getConnectionRequestTimeout())
        );

        // 设置 HTTP 客户端配置
        builder.setHttpClientConfigCallback(httpClientBuilder -> {
            // 设置连接池
            httpClientBuilder.setMaxConnTotal(properties.getMaxConnTotal());
            httpClientBuilder.setMaxConnPerRoute(properties.getMaxConnPerRoute());

            // 设置认证
            if (StringUtils.hasText(properties.getUsername())
                    && StringUtils.hasText(properties.getPassword())) {
                CredentialsProvider credentialsProvider = new BasicCredentialsProvider();
                credentialsProvider.setCredentials(
                        AuthScope.ANY,
                        new UsernamePasswordCredentials(
                                properties.getUsername(),
                                properties.getPassword()
                        )
                );
                httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider);
            }

            return httpClientBuilder;
        });

        // 设置失败重试策略
        builder.setFailureListener(new RestClient.FailureListener() {
            @Override
            public void onFailure(org.elasticsearch.client.Node node) {
                log.warn("ES 节点 [{}] 连接失败", node.getHost());
            }
        });

        restHighLevelClient = new RestHighLevelClient(builder);

        // 启用嗅探器
        if (properties.getEnableSniffer()) {
            sniffer = Sniffer.builder(restHighLevelClient.getLowLevelClient())
                    .setSniffIntervalMillis(properties.getSnifferInterval().intValue())
                    .build();
            log.info("ES 嗅探器已启用,间隔: {}ms", properties.getSnifferInterval());
        }

        log.info("ES 客户端初始化成功, 节点: {}", properties.getNodes());
        return restHighLevelClient;
    }

    @Bean
    @ConditionalOnMissingBean
    public ElasticsearchService elasticsearchService(
            RestHighLevelClient client,
            ElasticsearchProperties properties) {
        return new ElasticsearchService(client, properties);
    }

    @Bean
    @ConditionalOnMissingBean
    public ElasticsearchIndexService elasticsearchIndexService(
            RestHighLevelClient client,
            ElasticsearchProperties properties) {
        return new ElasticsearchIndexService(client, properties);
    }

    @Bean
    @ConditionalOnMissingBean
    public ElasticsearchDocumentService elasticsearchDocumentService(
            RestHighLevelClient client,
            ElasticsearchProperties properties) {
        return new ElasticsearchDocumentService(client, properties);
    }

    @PreDestroy
    public void destroy() {
        try {
            if (sniffer != null) {
                sniffer.close();
                log.info("ES 嗅探器已关闭");
            }
            if (restHighLevelClient != null) {
                restHighLevelClient.close();
                log.info("ES 客户端已关闭");
            }
        } catch (Exception e) {
            log.error("关闭 ES 客户端失败", e);
        }
    }
}

9.5 自定义异常类

package com.example.elasticsearch.exception;

import lombok.Getter;

/**
 * Elasticsearch 自定义异常
 */
@Getter
public class ElasticsearchException extends RuntimeException {

    private static final long serialVersionUID = 1L;

    /**
     * 索引名称
     */
    private String index;

    /**
     * 操作类型
     */
    private String operation;

    /**
     * 错误码
     */
    private String errorCode;

    /**
     * 是否可重试
     */
    private boolean retryable;

    public ElasticsearchException(String message) {
        super(message);
        this.retryable = false;
    }

    public ElasticsearchException(String message, Throwable cause) {
        super(message, cause);
        this.retryable = isRetryableException(cause);
    }

    public ElasticsearchException(String operation, String index, String message) {
        super(String.format("[%s] 索引 [%s] 操作失败: %s", operation, index, message));
        this.operation = operation;
        this.index = index;
        this.errorCode = operation + "_ERROR";
    }

    public ElasticsearchException(String operation, String index, Throwable cause) {
        super(String.format("[%s] 索引 [%s] 操作失败: %s",
                operation, index, cause.getMessage()), cause);
        this.operation = operation;
        this.index = index;
        this.errorCode = operation + "_ERROR";
        this.retryable = isRetryableException(cause);
    }

    /**
     * 判断异常是否可重试
     */
    private boolean isRetryableException(Throwable cause) {
        if (cause == null) {
            return false;
        }
        String message = cause.getMessage();
        if (message == null) {
            return false;
        }
        // 可重试的异常类型
        return message.contains("Connection refused")
                || message.contains("Connection reset")
                || message.contains("Connection timed out")
                || message.contains("Read timed out")
                || message.contains("No route to host")
                || message.contains("Service Unavailable")
                || message.contains("circuit_breaking_exception");
    }

    /**
     * 创建索引不存在异常
     */
    public static ElasticsearchException indexNotFound(String index) {
        ElasticsearchException ex = new ElasticsearchException(
                "INDEX_NOT_FOUND", index, "索引不存在");
        ex.errorCode = "INDEX_NOT_FOUND";
        return ex;
    }

    /**
     * 创建文档不存在异常
     */
    public static ElasticsearchException documentNotFound(String index, String id) {
        ElasticsearchException ex = new ElasticsearchException(
                "DOCUMENT_NOT_FOUND", index,
                String.format("文档 [%s] 不存在", id));
        ex.errorCode = "DOCUMENT_NOT_FOUND";
        return ex;
    }

    /**
     * 创建参数校验异常
     */
    public static ElasticsearchException invalidParameter(String message) {
        ElasticsearchException ex = new ElasticsearchException(message);
        ex.errorCode = "INVALID_PARAMETER";
        return ex;
    }
}

9.6 结果封装类

package com.example.elasticsearch.result;

import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.util.Collections;
import java.util.List;

/**
 * 分页结果封装
 */
@Data
@NoArgsConstructor
@AllArgsConstructor
public class PageResult<T> {

    /**
     * 数据列表
     */
    private List<T> records;

    /**
     * 总记录数
     */
    private Long total;

    /**
     * 当前页码
     */
    private Integer pageNum;

    /**
     * 每页大小
     */
    private Integer pageSize;

    /**
     * 总页数
     */
    private Integer totalPages;

    /**
     * 是否有下一页
     */
    private Boolean hasNext;

    /**
     * 是否有上一页
     */
    private Boolean hasPrevious;

    public PageResult(List<T> records, Long total, Integer pageNum, Integer pageSize) {
        this.records = records;
        this.total = total;
        this.pageNum = pageNum;
        this.pageSize = pageSize;
        this.totalPages = (int) Math.ceil((double) total / pageSize);
        this.hasNext = pageNum < totalPages;
        this.hasPrevious = pageNum > 1;
    }

    /**
     * 空结果
     */
    public static <T> PageResult<T> empty(Integer pageNum, Integer pageSize) {
        return new PageResult<>(Collections.emptyList(), 0L, pageNum, pageSize);
    }
}
package com.example.elasticsearch.result;

import lombok.Data;
import java.util.List;
import java.util.Map;

/**
 * 搜索结果封装(包含高亮、评分等信息)
 */
@Data
public class SearchResult<T> {

    /**
     * 文档 ID
     */
    private String documentId;

    /**
     * 数据对象
     */
    private T source;

    /**
     * 评分
     */
    private Float score;

    /**
     * 高亮字段
     */
    private Map<String, List<String>> highlight;

    /**
     * 排序值(用于深度分页)
     */
    private Object[] sortValues;

    public SearchResult(String documentId, T source) {
        this.documentId = documentId;
        this.source = source;
    }

    public SearchResult(String documentId, T source, Float score,
                        Map<String, List<String>> highlight) {
        this.documentId = documentId;
        this.source = source;
        this.score = score;
        this.highlight = highlight;
    }
}
package com.example.elasticsearch.result;

import lombok.Data;
import java.util.List;
import java.util.Map;

/**
 * 聚合结果封装
 */
@Data
public class AggregationResult {

    /**
     * 聚合名称
     */
    private String name;

    /**
     * 桶数据(terms 聚合)
     */
    private List<BucketData> buckets;

    /**
     * 数值(sum、avg、max、min 等)
     */
    private Double value;

    @Data
    public static class BucketData {
        private String key;
        private Long docCount;
        private Map<String, Object> subAggregations;
    }
}

9.7 查询构建器

package com.example.elasticsearch.query;

import lombok.Data;
import org.elasticsearch.index.query.*;
import org.elasticsearch.search.aggregations.AggregationBuilder;
import org.elasticsearch.search.aggregations.AggregationBuilders;
import org.elasticsearch.search.builder.SearchSourceBuilder;
import org.elasticsearch.search.fetch.subphase.highlight.HighlightBuilder;
import org.elasticsearch.search.sort.SortBuilder;
import org.elasticsearch.search.sort.SortBuilders;
import org.elasticsearch.search.sort.SortOrder;
import org.springframework.util.CollectionUtils;
import org.springframework.util.StringUtils;

import java.util.*;

/**
 * ES 查询构建器
 * 使用链式调用构建复杂查询
 */
@Data
public class EsQueryBuilder {

    /**
     * 索引名称
     */
    private String indexName;

    /**
     * Bool 查询条件
     */
    private BoolQueryBuilder boolQuery;

    /**
     * 排序条件
     */
    private List<SortBuilder<?>> sorts;

    /**
     * 高亮配置
     */
    private HighlightBuilder highlightBuilder;

    /**
     * 聚合配置
     */
    private List<AggregationBuilder> aggregations;

    /**
     * 分页参数
     */
    private Integer from;
    private Integer size;

    /**
     * 返回字段(包含)
     */
    private String[] includes;

    /**
     * 排除字段
     */
    private String[] excludes;

    /**
     * Search After(深度分页)
     */
    private Object[] searchAfter;

    /**
     * 是否追踪总数
     */
    private Boolean trackTotalHits = true;

    private EsQueryBuilder() {
        this.boolQuery = QueryBuilders.boolQuery();
        this.sorts = new ArrayList<>();
        this.aggregations = new ArrayList<>();
    }

    /**
     * 创建构建器
     */
    public static EsQueryBuilder builder(String indexName) {
        EsQueryBuilder builder = new EsQueryBuilder();
        builder.indexName = indexName;
        return builder;
    }

    // ==================== Must 条件(必须匹配)====================

    /**
     * 精确匹配(term)
     */
    public EsQueryBuilder term(String field, Object value) {
        if (value != null) {
            boolQuery.must(QueryBuilders.termQuery(field, value));
        }
        return this;
    }

    /**
     * 多值匹配(terms)
     */
    public EsQueryBuilder terms(String field, Collection<?> values) {
        if (!CollectionUtils.isEmpty(values)) {
            boolQuery.must(QueryBuilders.termsQuery(field, values));
        }
        return this;
    }

    /**
     * 全文匹配(match)
     */
    public EsQueryBuilder match(String field, Object value) {
        if (value != null && StringUtils.hasText(value.toString())) {
            boolQuery.must(QueryBuilders.matchQuery(field, value));
        }
        return this;
    }

    /**
     * 短语匹配(match_phrase)
     */
    public EsQueryBuilder matchPhrase(String field, Object value) {
        if (value != null && StringUtils.hasText(value.toString())) {
            boolQuery.must(QueryBuilders.matchPhraseQuery(field, value));
        }
        return this;
    }

    /**
     * 多字段匹配(multi_match)
     */
    public EsQueryBuilder multiMatch(Object value, String... fields) {
        if (value != null && StringUtils.hasText(value.toString()) && fields.length > 0) {
            boolQuery.must(QueryBuilders.multiMatchQuery(value, fields));
        }
        return this;
    }

    /**
     * 前缀匹配(prefix)
     */
    public EsQueryBuilder prefix(String field, String prefix) {
        if (StringUtils.hasText(prefix)) {
            boolQuery.must(QueryBuilders.prefixQuery(field, prefix));
        }
        return this;
    }

    /**
     * 通配符匹配(wildcard)
     */
    public EsQueryBuilder wildcard(String field, String pattern) {
        if (StringUtils.hasText(pattern)) {
            boolQuery.must(QueryBuilders.wildcardQuery(field, pattern));
        }
        return this;
    }

    /**
     * 范围查询
     */
    public EsQueryBuilder range(String field, Object gte, Object lte) {
        RangeQueryBuilder rangeQuery = QueryBuilders.rangeQuery(field);
        if (gte != null) {
            rangeQuery.gte(gte);
        }
        if (lte != null) {
            rangeQuery.lte(lte);
        }
        if (gte != null || lte != null) {
            boolQuery.must(rangeQuery);
        }
        return this;
    }

    /**
     * 范围查询(大于)
     */
    public EsQueryBuilder gt(String field, Object value) {
        if (value != null) {
            boolQuery.must(QueryBuilders.rangeQuery(field).gt(value));
        }
        return this;
    }

    /**
     * 范围查询(大于等于)
     */
    public EsQueryBuilder gte(String field, Object value) {
        if (value != null) {
            boolQuery.must(QueryBuilders.rangeQuery(field).gte(value));
        }
        return this;
    }

    /**
     * 范围查询(小于)
     */
    public EsQueryBuilder lt(String field, Object value) {
        if (value != null) {
            boolQuery.must(QueryBuilders.rangeQuery(field).lt(value));
        }
        return this;
    }

    /**
     * 范围查询(小于等于)
     */
    public EsQueryBuilder lte(String field, Object value) {
        if (value != null) {
            boolQuery.must(QueryBuilders.rangeQuery(field).lte(value));
        }
        return this;
    }

    /**
     * 存在字段查询
     */
    public EsQueryBuilder exists(String field) {
        boolQuery.must(QueryBuilders.existsQuery(field));
        return this;
    }

    // ==================== Filter 条件(过滤,不计算评分)====================

    /**
     * Filter - 精确匹配
     */
    public EsQueryBuilder filterTerm(String field, Object value) {
        if (value != null) {
            boolQuery.filter(QueryBuilders.termQuery(field, value));
        }
        return this;
    }

    /**
     * Filter - 多值匹配
     */
    public EsQueryBuilder filterTerms(String field, Collection<?> values) {
        if (!CollectionUtils.isEmpty(values)) {
            boolQuery.filter(QueryBuilders.termsQuery(field, values));
        }
        return this;
    }

    /**
     * Filter - 范围查询
     */
    public EsQueryBuilder filterRange(String field, Object gte, Object lte) {
        RangeQueryBuilder rangeQuery = QueryBuilders.rangeQuery(field);
        if (gte != null) {
            rangeQuery.gte(gte);
        }
        if (lte != null) {
            rangeQuery.lte(lte);
        }
        if (gte != null || lte != null) {
            boolQuery.filter(rangeQuery);
        }
        return this;
    }

    // ==================== Should 条件(或条件)====================

    /**
     * Should - 至少匹配一个
     */
    public EsQueryBuilder should(QueryBuilder... queries) {
        for (QueryBuilder query : queries) {
            boolQuery.should(query);
        }
        return this;
    }

    /**
     * Should - 多字段或查询
     */
    public EsQueryBuilder shouldMatch(String value, String... fields) {
        if (StringUtils.hasText(value) && fields.length > 0) {
            for (String field : fields) {
                boolQuery.should(QueryBuilders.matchQuery(field, value));
            }
            boolQuery.minimumShouldMatch(1);
        }
        return this;
    }

    /**
     * 设置最小 should 匹配数
     */
    public EsQueryBuilder minimumShouldMatch(int count) {
        boolQuery.minimumShouldMatch(count);
        return this;
    }

    // ==================== MustNot 条件(必须不匹配)====================

    /**
     * MustNot - 精确匹配
     */
    public EsQueryBuilder mustNotTerm(String field, Object value) {
        if (value != null) {
            boolQuery.mustNot(QueryBuilders.termQuery(field, value));
        }
        return this;
    }

    /**
     * MustNot - 多值匹配
     */
    public EsQueryBuilder mustNotTerms(String field, Collection<?> values) {
        if (!CollectionUtils.isEmpty(values)) {
            boolQuery.mustNot(QueryBuilders.termsQuery(field, values));
        }
        return this;
    }

    // ==================== 嵌套查询 ====================

    /**
     * 嵌套查询
     */
    public EsQueryBuilder nested(String path, QueryBuilder query) {
        boolQuery.must(QueryBuilders.nestedQuery(path, query,
                org.apache.lucene.search.join.ScoreMode.Avg));
        return this;
    }

    /**
     * 添加自定义查询条件
     */
    public EsQueryBuilder must(QueryBuilder query) {
        boolQuery.must(query);
        return this;
    }

    /**
     * 添加自定义过滤条件
     */
    public EsQueryBuilder filter(QueryBuilder query) {
        boolQuery.filter(query);
        return this;
    }

    // ==================== 排序 ====================

    /**
     * 添加排序
     */
    public EsQueryBuilder sort(String field, SortOrder order) {
        sorts.add(SortBuilders.fieldSort(field).order(order));
        return this;
    }

    /**
     * 按评分排序
     */
    public EsQueryBuilder sortByScore(SortOrder order) {
        sorts.add(SortBuilders.scoreSort().order(order));
        return this;
    }

    /**
     * 多字段排序
     */
    public EsQueryBuilder sorts(Map<String, SortOrder> sortMap) {
        sortMap.forEach((field, order) ->
                sorts.add(SortBuilders.fieldSort(field).order(order)));
        return this;
    }

    // ==================== 分页 ====================

    /**
     * 设置分页
     */
    public EsQueryBuilder page(int pageNum, int pageSize) {
        this.from = (pageNum - 1) * pageSize;
        this.size = pageSize;
        return this;
    }

    /**
     * 设置起始位置和大小
     */
    public EsQueryBuilder fromSize(int from, int size) {
        this.from = from;
        this.size = size;
        return this;
    }

    /**
     * 深度分页(Search After)
     */
    public EsQueryBuilder searchAfter(Object[] values) {
        this.searchAfter = values;
        return this;
    }

    // ==================== 高亮 ====================

    /**
     * 添加高亮字段
     */
    public EsQueryBuilder highlight(String... fields) {
        if (fields.length > 0) {
            highlightBuilder = new HighlightBuilder();
            for (String field : fields) {
                highlightBuilder.field(field);
            }
            highlightBuilder.preTags("<em class='highlight'>");
            highlightBuilder.postTags("</em>");
        }
        return this;
    }

    /**
     * 自定义高亮配置
     */
    public EsQueryBuilder highlight(String preTag, String postTag, String... fields) {
        if (fields.length > 0) {
            highlightBuilder = new HighlightBuilder();
            for (String field : fields) {
                highlightBuilder.field(field);
            }
            highlightBuilder.preTags(preTag);
            highlightBuilder.postTags(postTag);
        }
        return this;
    }

    // ==================== 聚合 ====================

    /**
     * Terms 聚合
     */
    public EsQueryBuilder termsAggregation(String name, String field, int size) {
        aggregations.add(AggregationBuilders.terms(name).field(field).size(size));
        return this;
    }

    /**
     * Sum 聚合
     */
    public EsQueryBuilder sumAggregation(String name, String field) {
        aggregations.add(AggregationBuilders.sum(name).field(field));
        return this;
    }

    /**
     * Avg 聚合
     */
    public EsQueryBuilder avgAggregation(String name, String field) {
        aggregations.add(AggregationBuilders.avg(name).field(field));
        return this;
    }

    /**
     * Max 聚合
     */
    public EsQueryBuilder maxAggregation(String name, String field) {
        aggregations.add(AggregationBuilders.max(name).field(field));
        return this;
    }

    /**
     * Min 聚合
     */
    public EsQueryBuilder minAggregation(String name, String field) {
        aggregations.add(AggregationBuilders.min(name).field(field));
        return this;
    }

    /**
     * 日期直方图聚合
     */
    public EsQueryBuilder dateHistogramAggregation(String name, String field,
                                                    String interval) {
        aggregations.add(AggregationBuilders.dateHistogram(name)
                .field(field)
                .calendarInterval(new org.elasticsearch.search.aggregations.bucket
                        .histogram.DateHistogramInterval(interval)));
        return this;
    }

    /**
     * 添加自定义聚合
     */
    public EsQueryBuilder aggregation(AggregationBuilder aggregation) {
        aggregations.add(aggregation);
        return this;
    }

    // ==================== 返回字段 ====================

    /**
     * 指定返回字段
     */
    public EsQueryBuilder includes(String... fields) {
        this.includes = fields;
        return this;
    }

    /**
     * 排除字段
     */
    public EsQueryBuilder excludes(String... fields) {
        this.excludes = fields;
        return this;
    }

    /**
     * 是否追踪总数
     */
    public EsQueryBuilder trackTotalHits(boolean track) {
        this.trackTotalHits = track;
        return this;
    }

    // ==================== 构建 ====================

    /**
     * 构建 SearchSourceBuilder
     */
    public SearchSourceBuilder build() {
        SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();

        // 查询条件
        sourceBuilder.query(boolQuery);

        // 分页
        if (from != null) {
            sourceBuilder.from(from);
        }
        if (size != null) {
            sourceBuilder.size(size);
        }

        // 排序
        for (SortBuilder<?> sort : sorts) {
            sourceBuilder.sort(sort);
        }

        // 高亮
        if (highlightBuilder != null) {
            sourceBuilder.highlighter(highlightBuilder);
        }

        // 聚合
        for (AggregationBuilder aggregation : aggregations) {
            sourceBuilder.aggregation(aggregation);
        }

        // 返回字段
        if (includes != null || excludes != null) {
            sourceBuilder.fetchSource(includes, excludes);
        }

        // Search After
        if (searchAfter != null) {
            sourceBuilder.searchAfter(searchAfter);
        }

        // 追踪总数
        sourceBuilder.trackTotalHits(trackTotalHits);

        return sourceBuilder;
    }
}

9.8 重试工具类

package com.example.elasticsearch.util;

import com.example.elasticsearch.exception.ElasticsearchException;
import lombok.extern.slf4j.Slf4j;

import java.util.concurrent.Callable;
import java.util.function.Predicate;

/**
 * 重试工具类
 */
@Slf4j
public class RetryUtil {

    /**
     * 执行带重试的操作
     *
     * @param callable      要执行的操作
     * @param maxRetries    最大重试次数
     * @param retryInterval 重试间隔(毫秒)
     * @param retryOn       判断是否需要重试的条件
     * @param operationName 操作名称(用于日志)
     * @return 操作结果
     */
    public static <T> T executeWithRetry(
            Callable<T> callable,
            int maxRetries,
            long retryInterval,
            Predicate<Exception> retryOn,
            String operationName) {

        Exception lastException = null;

        for (int attempt = 0; attempt <= maxRetries; attempt++) {
            try {
                return callable.call();
            } catch (Exception e) {
                lastException = e;

                // 判断是否需要重试
                if (attempt < maxRetries && retryOn.test(e)) {
                    log.warn("[{}] 操作失败,第 {}/{} 次重试,错误: {}",
                            operationName, attempt + 1, maxRetries, e.getMessage());

                    try {
                        Thread.sleep(retryInterval * (attempt + 1)); // 指数退避
                    } catch (InterruptedException ie) {
                        Thread.currentThread().interrupt();
                        throw new ElasticsearchException("重试被中断", ie);
                    }
                } else {
                    break;
                }
            }
        }

        log.error("[{}] 操作失败,已达最大重试次数", operationName, lastException);
        if (lastException instanceof ElasticsearchException) {
            throw (ElasticsearchException) lastException;
        }
        throw new ElasticsearchException(operationName + " 操作失败", lastException);
    }

    /**
     * 执行带重试的操作(无返回值)
     */
    public static void executeWithRetry(
            Runnable runnable,
            int maxRetries,
            long retryInterval,
            Predicate<Exception> retryOn,
            String operationName) {

        executeWithRetry(() -> {
            runnable.run();
            return null;
        }, maxRetries, retryInterval, retryOn, operationName);
    }

    /**
     * 默认的重试条件判断
     */
    public static Predicate<Exception> defaultRetryCondition() {
        return e -> {
            if (e instanceof ElasticsearchException) {
                return ((ElasticsearchException) e).isRetryable();
            }
            String message = e.getMessage();
            if (message == null) {
                return false;
            }
            return message.contains("Connection")
                    || message.contains("timed out")
                    || message.contains("Unavailable");
        };
    }
}

9.9 参数校验工具类

package com.example.elasticsearch.util;

import com.example.elasticsearch.exception.ElasticsearchException;
import org.springframework.util.CollectionUtils;
import org.springframework.util.StringUtils;

import java.util.Collection;

/**
 * 参数校验工具类
 */
public class ParamValidator {

    private ParamValidator() {}

    /**
     * 校验索引名称
     */
    public static void validateIndexName(String indexName) {
        if (!StringUtils.hasText(indexName)) {
            throw ElasticsearchException.invalidParameter("索引名称不能为空");
        }
        if (indexName.contains(" ")) {
            throw ElasticsearchException.invalidParameter("索引名称不能包含空格");
        }
        if (!indexName.equals(indexName.toLowerCase())) {
            throw ElasticsearchException.invalidParameter("索引名称必须小写");
        }
    }

    /**
     * 校验文档ID
     */
    public static void validateDocumentId(String documentId) {
        if (!StringUtils.hasText(documentId)) {
            throw ElasticsearchException.invalidParameter("文档ID不能为空");
        }
    }

    /**
     * 校验分页参数
     */
    public static void validatePageParam(int pageNum, int pageSize) {
        if (pageNum < 1) {
            throw ElasticsearchException.invalidParameter("页码必须大于0");
        }
        if (pageSize < 1 || pageSize > 10000) {
            throw ElasticsearchException.invalidParameter("每页大小必须在1-10000之间");
        }
        // ES 默认限制 from + size <= 10000
        if ((long) (pageNum - 1) * pageSize + pageSize > 10000) {
            throw ElasticsearchException.invalidParameter(
                    "分页深度超出限制,请使用 searchAfter 方式");
        }
    }

    /**
     * 校验批量数据
     */
    public static void validateBatchData(Collection<?> dataList) {
        if (CollectionUtils.isEmpty(dataList)) {
            throw ElasticsearchException.invalidParameter("批量数据不能为空");
        }
    }

    /**
     * 校验非空
     */
    public static void notNull(Object object, String message) {
        if (object == null) {
            throw ElasticsearchException.invalidParameter(message);
        }
    }

    /**
     * 校验字符串非空
     */
    public static void notBlank(String str, String message) {
        if (!StringUtils.hasText(str)) {
            throw ElasticsearchException.invalidParameter(message);
        }
    }
}

9.10 核心 Service 封装

package com.example.elasticsearch.core;

import com.alibaba.fastjson.JSON;
import com.example.elasticsearch.config.ElasticsearchProperties;
import com.example.elasticsearch.exception.ElasticsearchException;
import com.example.elasticsearch.query.EsQueryBuilder;
import com.example.elasticsearch.result.AggregationResult;
import com.example.elasticsearch.result.PageResult;
import com.example.elasticsearch.result.ScrollResult;
import com.example.elasticsearch.result.SearchResult;
import com.example.elasticsearch.util.ParamValidator;
import com.example.elasticsearch.util.RetryUtil;
import lombok.extern.slf4j.Slf4j;
import org.elasticsearch.action.search.*;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.common.text.Text;
import org.elasticsearch.common.unit.TimeValue;
import org.elasticsearch.index.query.QueryBuilder;
import org.elasticsearch.search.SearchHit;
import org.elasticsearch.search.SearchHits;
import org.elasticsearch.search.aggregations.Aggregation;
import org.elasticsearch.search.aggregations.bucket.terms.Terms;
import org.elasticsearch.search.aggregations.metrics.*;
import org.elasticsearch.search.builder.SearchSourceBuilder;
import org.elasticsearch.search.fetch.subphase.highlight.HighlightField;

import java.io.IOException;
import java.util.*;
import java.util.stream.Collectors;

/**
 * Elasticsearch 核心搜索服务
 *
 * 特性:
 * - 完善的参数校验
 * - 自动重试机制
 * - 慢查询日志
 * - 滚动查询支持
 * - 空值安全处理
 */
@Slf4j
public class ElasticsearchService {

    private final RestHighLevelClient client;
    private final ElasticsearchProperties properties;

    public ElasticsearchService(RestHighLevelClient client,
                                ElasticsearchProperties properties) {
        this.client = client;
        this.properties = properties;
    }

    // ==================== 基础查询 ====================

    /**
     * 简单查询 - 根据单个字段精确匹配
     *
     * @param indexName 索引名称
     * @param field     字段名
     * @param value     字段值
     * @param clazz     返回类型
     * @return 匹配的文档列表
     */
    public <T> List<T> query(String indexName, String field, Object value,
                             Class<T> clazz) {
        ParamValidator.validateIndexName(indexName);
        ParamValidator.notBlank(field, "查询字段不能为空");

        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(indexName)
                .term(field, value);
        return search(queryBuilder, clazz);
    }

    /**
     * 多条件查询 - 根据多个字段匹配
     *
     * @param indexName 索引名称
     * @param params    查询参数 (字段名 -> 字段值)
     * @param clazz     返回类型
     * @return 匹配的文档列表
     */
    public <T> List<T> query(String indexName, Map<String, Object> params,
                             Class<T> clazz) {
        ParamValidator.validateIndexName(indexName);

        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(indexName);

        if (params != null && !params.isEmpty()) {
            params.forEach((field, value) -> {
                if (value != null) {
                    if (value instanceof String && ((String) value).length() > 0) {
                        queryBuilder.match(field, value);
                    } else if (!(value instanceof String)) {
                        queryBuilder.filterTerm(field, value);
                    }
                }
            });
        }

        return search(queryBuilder, clazz);
    }

    /**
     * 分页查询
     *
     * @param indexName 索引名称
     * @param params    查询参数
     * @param pageNum   页码(从1开始)
     * @param pageSize  每页大小
     * @param clazz     返回类型
     * @return 分页结果
     */
    public <T> PageResult<T> queryPage(String indexName, Map<String, Object> params,
                                       int pageNum, int pageSize, Class<T> clazz) {
        ParamValidator.validateIndexName(indexName);
        ParamValidator.validatePageParam(pageNum, pageSize);

        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(indexName)
                .page(pageNum, pageSize);

        if (params != null && !params.isEmpty()) {
            params.forEach((field, value) -> {
                if (value != null) {
                    if (value instanceof String && ((String) value).length() > 0) {
                        queryBuilder.match(field, value);
                    } else if (!(value instanceof String)) {
                        queryBuilder.filterTerm(field, value);
                    }
                }
            });
        }

        return searchPage(queryBuilder, pageNum, pageSize, clazz);
    }

    // ==================== 高级查询 ====================

    /**
     * 使用查询构建器执行查询
     *
     * @param queryBuilder 查询构建器
     * @param clazz        返回类型
     * @return 匹配的文档列表
     */
    public <T> List<T> search(EsQueryBuilder queryBuilder, Class<T> clazz) {
        ParamValidator.validateIndexName(queryBuilder.getIndexName());
        ParamValidator.notNull(clazz, "返回类型不能为空");

        return RetryUtil.executeWithRetry(
                () -> doSearch(queryBuilder, clazz),
                properties.getRetryTimes(),
                properties.getRetryInterval(),
                RetryUtil.defaultRetryCondition(),
                "SEARCH"
        );
    }

    private <T> List<T> doSearch(EsQueryBuilder queryBuilder, Class<T> clazz)
            throws IOException {
        long startTime = System.currentTimeMillis();

        try {
            SearchSourceBuilder sourceBuilder = queryBuilder.build();
            SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
            searchRequest.source(sourceBuilder);

            // 打印 DSL
            if (properties.getPrintDsl()) {
                log.info("ES 查询 DSL: {}", sourceBuilder.toString());
            }

            SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);

            // 检查响应状态
            checkResponseStatus(response, queryBuilder.getIndexName());

            return parseHits(response.getHits(), clazz);
        } finally {
            logSlowQuery(startTime, "SEARCH", queryBuilder.getIndexName());
        }
    }

    /**
     * 分页查询
     */
    public <T> PageResult<T> searchPage(EsQueryBuilder queryBuilder,
                                         int pageNum, int pageSize, Class<T> clazz) {
        ParamValidator.validateIndexName(queryBuilder.getIndexName());
        ParamValidator.validatePageParam(pageNum, pageSize);
        ParamValidator.notNull(clazz, "返回类型不能为空");

        return RetryUtil.executeWithRetry(
                () -> doSearchPage(queryBuilder, pageNum, pageSize, clazz),
                properties.getRetryTimes(),
                properties.getRetryInterval(),
                RetryUtil.defaultRetryCondition(),
                "SEARCH_PAGE"
        );
    }

    private <T> PageResult<T> doSearchPage(EsQueryBuilder queryBuilder,
                                            int pageNum, int pageSize,
                                            Class<T> clazz) throws IOException {
        long startTime = System.currentTimeMillis();

        try {
            queryBuilder.page(pageNum, pageSize);

            SearchSourceBuilder sourceBuilder = queryBuilder.build();
            SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
            searchRequest.source(sourceBuilder);

            if (properties.getPrintDsl()) {
                log.info("ES 分页查询 DSL: {}", sourceBuilder.toString());
            }

            SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);

            checkResponseStatus(response, queryBuilder.getIndexName());

            List<T> records = parseHits(response.getHits(), clazz);
            long total = getTotalHits(response.getHits());

            return new PageResult<>(records, total, pageNum, pageSize);
        } finally {
            logSlowQuery(startTime, "SEARCH_PAGE", queryBuilder.getIndexName());
        }
    }

    /**
     * 带高亮的查询
     */
    public <T> List<SearchResult<T>> searchWithHighlight(EsQueryBuilder queryBuilder,
                                                          Class<T> clazz) {
        ParamValidator.validateIndexName(queryBuilder.getIndexName());

        return RetryUtil.executeWithRetry(
                () -> doSearchWithHighlight(queryBuilder, clazz),
                properties.getRetryTimes(),
                properties.getRetryInterval(),
                RetryUtil.defaultRetryCondition(),
                "SEARCH_HIGHLIGHT"
        );
    }

    private <T> List<SearchResult<T>> doSearchWithHighlight(
            EsQueryBuilder queryBuilder, Class<T> clazz) throws IOException {
        long startTime = System.currentTimeMillis();

        try {
            SearchSourceBuilder sourceBuilder = queryBuilder.build();
            SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
            searchRequest.source(sourceBuilder);

            if (properties.getPrintDsl()) {
                log.info("ES 高亮查询 DSL: {}", sourceBuilder.toString());
            }

            SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);

            checkResponseStatus(response, queryBuilder.getIndexName());

            return parseHitsWithHighlight(response.getHits(), clazz);
        } finally {
            logSlowQuery(startTime, "SEARCH_HIGHLIGHT", queryBuilder.getIndexName());
        }
    }

    /**
     * 带高亮的分页查询
     */
    public <T> PageResult<SearchResult<T>> searchPageWithHighlight(
            EsQueryBuilder queryBuilder, int pageNum, int pageSize, Class<T> clazz) {
        ParamValidator.validateIndexName(queryBuilder.getIndexName());
        ParamValidator.validatePageParam(pageNum, pageSize);

        return RetryUtil.executeWithRetry(
                () -> doSearchPageWithHighlight(queryBuilder, pageNum, pageSize, clazz),
                properties.getRetryTimes(),
                properties.getRetryInterval(),
                RetryUtil.defaultRetryCondition(),
                "SEARCH_PAGE_HIGHLIGHT"
        );
    }

    private <T> PageResult<SearchResult<T>> doSearchPageWithHighlight(
            EsQueryBuilder queryBuilder, int pageNum, int pageSize,
            Class<T> clazz) throws IOException {
        long startTime = System.currentTimeMillis();

        try {
            queryBuilder.page(pageNum, pageSize);

            SearchSourceBuilder sourceBuilder = queryBuilder.build();
            SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
            searchRequest.source(sourceBuilder);

            SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);

            checkResponseStatus(response, queryBuilder.getIndexName());

            List<SearchResult<T>> records = parseHitsWithHighlight(response.getHits(), clazz);
            long total = getTotalHits(response.getHits());

            return new PageResult<>(records, total, pageNum, pageSize);
        } finally {
            logSlowQuery(startTime, "SEARCH_PAGE_HIGHLIGHT", queryBuilder.getIndexName());
        }
    }

    // ==================== 滚动查询(大数据量) ====================

    /**
     * 滚动查询 - 初始化
     * 适用于导出大量数据的场景
     *
     * @param queryBuilder 查询构建器
     * @param scrollTime   滚动上下文保持时间(分钟)
     * @param size         每批大小
     * @param clazz        返回类型
     * @return 滚动结果(包含scrollId和首批数据)
     */
    public <T> ScrollResult<T> scrollSearch(EsQueryBuilder queryBuilder,
                                             int scrollTime, int size,
                                             Class<T> clazz) {
        ParamValidator.validateIndexName(queryBuilder.getIndexName());

        try {
            queryBuilder.fromSize(0, size);

            SearchSourceBuilder sourceBuilder = queryBuilder.build();
            SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
            searchRequest.source(sourceBuilder);
            searchRequest.scroll(TimeValue.timeValueMinutes(scrollTime));

            if (properties.getPrintDsl()) {
                log.info("ES 滚动查询初始化 DSL: {}", sourceBuilder.toString());
            }

            SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);

            checkResponseStatus(response, queryBuilder.getIndexName());

            List<T> records = parseHits(response.getHits(), clazz);
            long total = getTotalHits(response.getHits());
            String scrollId = response.getScrollId();

            return new ScrollResult<>(scrollId, records, total, records.size() < size);
        } catch (IOException e) {
            throw new ElasticsearchException("SCROLL_SEARCH",
                    queryBuilder.getIndexName(), e);
        }
    }

    /**
     * 滚动查询 - 继续获取下一批
     *
     * @param scrollId   滚动ID
     * @param scrollTime 滚动上下文保持时间(分钟)
     * @param size       每批大小
     * @param clazz      返回类型
     * @return 滚动结果
     */
    public <T> ScrollResult<T> scrollNext(String scrollId, int scrollTime,
                                           int size, Class<T> clazz) {
        ParamValidator.notBlank(scrollId, "scrollId不能为空");

        try {
            SearchScrollRequest scrollRequest = new SearchScrollRequest(scrollId);
            scrollRequest.scroll(TimeValue.timeValueMinutes(scrollTime));

            SearchResponse response = client.scroll(scrollRequest, RequestOptions.DEFAULT);

            List<T> records = parseHits(response.getHits(), clazz);
            String newScrollId = response.getScrollId();

            return new ScrollResult<>(newScrollId, records,
                    getTotalHits(response.getHits()), records.size() < size);
        } catch (IOException e) {
            throw new ElasticsearchException("滚动查询失败", e);
        }
    }

    /**
     * 清除滚动上下文
     *
     * @param scrollIds 滚动ID列表
     */
    public void clearScroll(String... scrollIds) {
        if (scrollIds == null || scrollIds.length == 0) {
            return;
        }

        try {
            ClearScrollRequest clearScrollRequest = new ClearScrollRequest();
            clearScrollRequest.scrollIds(Arrays.asList(scrollIds));

            ClearScrollResponse response = client.clearScroll(
                    clearScrollRequest, RequestOptions.DEFAULT);

            if (!response.isSucceeded()) {
                log.warn("清除滚动上下文失败");
            }
        } catch (IOException e) {
            log.warn("清除滚动上下文异常", e);
        }
    }

    // ==================== 聚合查询 ====================

    /**
     * 聚合查询
     */
    public Map<String, AggregationResult> searchAggregation(EsQueryBuilder queryBuilder) {
        ParamValidator.validateIndexName(queryBuilder.getIndexName());

        return RetryUtil.executeWithRetry(
                () -> doSearchAggregation(queryBuilder),
                properties.getRetryTimes(),
                properties.getRetryInterval(),
                RetryUtil.defaultRetryCondition(),
                "SEARCH_AGGREGATION"
        );
    }

    private Map<String, AggregationResult> doSearchAggregation(
            EsQueryBuilder queryBuilder) throws IOException {
        long startTime = System.currentTimeMillis();

        try {
            queryBuilder.fromSize(0, 0);

            SearchSourceBuilder sourceBuilder = queryBuilder.build();
            SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
            searchRequest.source(sourceBuilder);

            if (properties.getPrintDsl()) {
                log.info("ES 聚合查询 DSL: {}", sourceBuilder.toString());
            }

            SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);

            checkResponseStatus(response, queryBuilder.getIndexName());

            return parseAggregations(response.getAggregations());
        } finally {
            logSlowQuery(startTime, "SEARCH_AGGREGATION", queryBuilder.getIndexName());
        }
    }

    // ==================== Search After(深度分页) ====================

    /**
     * Search After 查询
     * 适用于深度分页场景,避免 from+size 的 10000 限制
     *
     * @param queryBuilder      查询构建器
     * @param searchAfterValues 上一页最后一条记录的排序值
     * @param size              每页大小
     * @param clazz             返回类型
     * @return 搜索结果列表
     */
    public <T> List<SearchResult<T>> searchAfter(EsQueryBuilder queryBuilder,
                                                   Object[] searchAfterValues,
                                                   int size, Class<T> clazz) {
        ParamValidator.validateIndexName(queryBuilder.getIndexName());

        return RetryUtil.executeWithRetry(
                () -> doSearchAfter(queryBuilder, searchAfterValues, size, clazz),
                properties.getRetryTimes(),
                properties.getRetryInterval(),
                RetryUtil.defaultRetryCondition(),
                "SEARCH_AFTER"
        );
    }

    private <T> List<SearchResult<T>> doSearchAfter(
            EsQueryBuilder queryBuilder, Object[] searchAfterValues,
            int size, Class<T> clazz) throws IOException {
        long startTime = System.currentTimeMillis();

        try {
            queryBuilder.fromSize(0, size);
            if (searchAfterValues != null && searchAfterValues.length > 0) {
                queryBuilder.searchAfter(searchAfterValues);
            }

            SearchSourceBuilder sourceBuilder = queryBuilder.build();
            SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
            searchRequest.source(sourceBuilder);

            if (properties.getPrintDsl()) {
                log.info("ES Search After DSL: {}", sourceBuilder.toString());
            }

            SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);

            checkResponseStatus(response, queryBuilder.getIndexName());

            return parseHitsWithHighlight(response.getHits(), clazz);
        } finally {
            logSlowQuery(startTime, "SEARCH_AFTER", queryBuilder.getIndexName());
        }
    }

    // ==================== 统计查询 ====================

    /**
     * 统计符合条件的文档数量
     */
    public long count(EsQueryBuilder queryBuilder) {
        ParamValidator.validateIndexName(queryBuilder.getIndexName());

        return RetryUtil.executeWithRetry(
                () -> doCount(queryBuilder),
                properties.getRetryTimes(),
                properties.getRetryInterval(),
                RetryUtil.defaultRetryCondition(),
                "COUNT"
        );
    }

    private long doCount(EsQueryBuilder queryBuilder) throws IOException {
        queryBuilder.fromSize(0, 0);

        SearchSourceBuilder sourceBuilder = queryBuilder.build();
        SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
        searchRequest.source(sourceBuilder);

        SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);

        checkResponseStatus(response, queryBuilder.getIndexName());

        return getTotalHits(response.getHits());
    }

    /**
     * 判断是否存在符合条件的文档
     */
    public boolean exists(EsQueryBuilder queryBuilder) {
        return count(queryBuilder) > 0;
    }

    // ==================== 原生查询 ====================

    /**
     * 执行原生 QueryBuilder 查询
     * 适用于封装方法无法满足的复杂场景
     */
    public <T> List<T> searchByQueryBuilder(String indexName,
                                             QueryBuilder queryBuilder,
                                             int size, Class<T> clazz) {
        ParamValidator.validateIndexName(indexName);
        ParamValidator.notNull(queryBuilder, "查询条件不能为空");

        try {
            SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
            sourceBuilder.query(queryBuilder);
            sourceBuilder.size(size);
            sourceBuilder.trackTotalHits(true);

            SearchRequest searchRequest = new SearchRequest(indexName);
            searchRequest.source(sourceBuilder);

            SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);

            checkResponseStatus(response, indexName);

            return parseHits(response.getHits(), clazz);
        } catch (IOException e) {
            throw new ElasticsearchException("SEARCH_BY_QUERY_BUILDER", indexName, e);
        }
    }

    // ==================== 内部工具方法 ====================

    /**
     * 解析命中结果
     */
    private <T> List<T> parseHits(SearchHits hits, Class<T> clazz) {
        if (hits == null || hits.getHits() == null) {
            return Collections.emptyList();
        }

        List<T> result = new ArrayList<>();

        for (SearchHit hit : hits.getHits()) {
            try {
                String source = hit.getSourceAsString();
                if (source != null && !source.isEmpty()) {
                    T obj = JSON.parseObject(source, clazz);
                    result.add(obj);
                }
            } catch (Exception e) {
                log.warn("解析文档失败, id: {}, error: {}", hit.getId(), e.getMessage());
            }
        }

        return result;
    }

    /**
     * 解析带高亮的命中结果
     */
    private <T> List<SearchResult<T>> parseHitsWithHighlight(SearchHits hits,
                                                              Class<T> clazz) {
        if (hits == null || hits.getHits() == null) {
            return Collections.emptyList();
        }

        List<SearchResult<T>> result = new ArrayList<>();

        for (SearchHit hit : hits.getHits()) {
            try {
                String source = hit.getSourceAsString();
                if (source == null || source.isEmpty()) {
                    continue;
                }

                T obj = JSON.parseObject(source, clazz);

                // 解析高亮
                Map<String, List<String>> highlightMap = new HashMap<>();
                Map<String, HighlightField> highlightFields = hit.getHighlightFields();

                if (highlightFields != null && !highlightFields.isEmpty()) {
                    highlightFields.forEach((field, highlightField) -> {
                        if (highlightField.getFragments() != null) {
                            List<String> fragments = Arrays.stream(highlightField.getFragments())
                                    .map(Text::string)
                                    .collect(Collectors.toList());
                            highlightMap.put(field, fragments);
                        }
                    });
                }

                SearchResult<T> searchResult = new SearchResult<>(
                        hit.getId(), obj, hit.getScore(), highlightMap);
                searchResult.setSortValues(hit.getSortValues());

                result.add(searchResult);
            } catch (Exception e) {
                log.warn("解析文档失败, id: {}, error: {}", hit.getId(), e.getMessage());
            }
        }

        return result;
    }

    /**
     * 解析聚合结果
     */
    private Map<String, AggregationResult> parseAggregations(
            org.elasticsearch.search.aggregations.Aggregations aggregations) {

        Map<String, AggregationResult> result = new HashMap<>();

        if (aggregations == null) {
            return result;
        }

        for (Aggregation aggregation : aggregations) {
            try {
                AggregationResult aggResult = new AggregationResult();
                aggResult.setName(aggregation.getName());

                if (aggregation instanceof Terms) {
                    Terms terms = (Terms) aggregation;
                    List<AggregationResult.BucketData> buckets = new ArrayList<>();

                    for (Terms.Bucket bucket : terms.getBuckets()) {
                        AggregationResult.BucketData bucketData =
                                new AggregationResult.BucketData();
                        bucketData.setKey(bucket.getKeyAsString());
                        bucketData.setDocCount(bucket.getDocCount());
                        buckets.add(bucketData);
                    }
                    aggResult.setBuckets(buckets);

                } else if (aggregation instanceof Sum) {
                    aggResult.setValue(((Sum) aggregation).getValue());
                } else if (aggregation instanceof Avg) {
                    aggResult.setValue(((Avg) aggregation).getValue());
                } else if (aggregation instanceof Max) {
                    aggResult.setValue(((Max) aggregation).getValue());
                } else if (aggregation instanceof Min) {
                    aggResult.setValue(((Min) aggregation).getValue());
                } else if (aggregation instanceof ValueCount) {
                    aggResult.setValue((double) ((ValueCount) aggregation).getValue());
                } else if (aggregation instanceof Cardinality) {
                    aggResult.setValue((double) ((Cardinality) aggregation).getValue());
                }

                result.put(aggregation.getName(), aggResult);
            } catch (Exception e) {
                log.warn("解析聚合结果失败, name: {}, error: {}",
                        aggregation.getName(), e.getMessage());
            }
        }

        return result;
    }

    /**
     * 获取总命中数(兼容不同版本)
     */
    private long getTotalHits(SearchHits hits) {
        if (hits == null || hits.getTotalHits() == null) {
            return 0L;
        }
        return hits.getTotalHits().value;
    }

    /**
     * 检查响应状态
     */
    private void checkResponseStatus(SearchResponse response, String indexName) {
        if (response == null) {
            throw new ElasticsearchException("SEARCH", indexName, "响应为空");
        }

        if (response.isTimedOut()) {
            log.warn("ES 查询超时, index: {}", indexName);
        }

        if (response.getShardFailures() != null && response.getShardFailures().length > 0) {
            log.warn("ES 查询部分分片失败, index: {}, failures: {}",
                    indexName, response.getShardFailures().length);
        }
    }

    /**
     * 记录慢查询日志
     */
    private void logSlowQuery(long startTime, String operation, String indexName) {
        long elapsed = System.currentTimeMillis() - startTime;

        if (elapsed >= properties.getSlowQueryThreshold()) {
            log.warn("[慢查询] 操作: {}, 索引: {}, 耗时: {}ms",
                    operation, indexName, elapsed);
        } else {
            log.debug("ES 操作完成, 操作: {}, 索引: {}, 耗时: {}ms",
                    operation, indexName, elapsed);
        }
    }
}

9.11 文档操作服务

package com.example.elasticsearch.core;

import com.alibaba.fastjson.JSON;
import com.example.elasticsearch.config.ElasticsearchProperties;
import com.example.elasticsearch.exception.ElasticsearchException;
import com.example.elasticsearch.util.ParamValidator;
import com.example.elasticsearch.util.RetryUtil;
import lombok.extern.slf4j.Slf4j;
import org.elasticsearch.action.DocWriteResponse;
import org.elasticsearch.action.bulk.BulkItemResponse;
import org.elasticsearch.action.bulk.BulkRequest;
import org.elasticsearch.action.bulk.BulkResponse;
import org.elasticsearch.action.delete.DeleteRequest;
import org.elasticsearch.action.delete.DeleteResponse;
import org.elasticsearch.action.get.*;
import org.elasticsearch.action.index.IndexRequest;
import org.elasticsearch.action.index.IndexResponse;
import org.elasticsearch.action.support.WriteRequest;
import org.elasticsearch.action.update.UpdateRequest;
import org.elasticsearch.action.update.UpdateResponse;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.common.xcontent.XContentType;
import org.elasticsearch.index.query.QueryBuilder;
import org.elasticsearch.index.reindex.BulkByScrollResponse;
import org.elasticsearch.index.reindex.DeleteByQueryRequest;
import org.elasticsearch.index.reindex.UpdateByQueryRequest;
import org.elasticsearch.script.Script;
import org.elasticsearch.script.ScriptType;
import org.springframework.util.CollectionUtils;
import org.springframework.util.StringUtils;

import java.io.IOException;
import java.util.*;

/**
 * Elasticsearch 文档操作服务(增强版)
 *
 * 特性:
 * - 完善的参数校验
 * - 批量操作分批处理
 * - 自动重试机制
 * - 详细的操作日志
 */
@Slf4j
public class ElasticsearchDocumentService {

    private final RestHighLevelClient client;
    private final ElasticsearchProperties properties;

    public ElasticsearchDocumentService(RestHighLevelClient client,
                                         ElasticsearchProperties properties) {
        this.client = client;
        this.properties = properties;
    }

    // ==================== 新增操作 ====================

    /**
     * 添加文档(自动生成ID)
     */
    public String add(String indexName, Object data) {
        return add(indexName, null, data);
    }

    /**
     * 添加文档(指定ID)
     */
    public String add(String indexName, String id, Object data) {
        ParamValidator.validateIndexName(indexName);
        ParamValidator.notNull(data, "文档数据不能为空");

        return RetryUtil.executeWithRetry(
                () -> doAdd(indexName, id, data),
                properties.getRetryTimes(),
                properties.getRetryInterval(),
                RetryUtil.defaultRetryCondition(),
                "ADD_DOCUMENT"
        );
    }

    private String doAdd(String indexName, String id, Object data) throws IOException {
        IndexRequest request = new IndexRequest(indexName);

        if (StringUtils.hasText(id)) {
            request.id(id);
        }

        String jsonData = JSON.toJSONString(data);
        request.source(jsonData, XContentType.JSON);

        // 兼容 ES 6.x
        if (properties.getEnableType()) {
            request.type(properties.getDefaultType());
        }

        // 设置刷新策略
        setRefreshPolicy(request);

        IndexResponse response = client.index(request, RequestOptions.DEFAULT);

        log.info("添加文档成功, index: {}, id: {}", indexName, response.getId());
        return response.getId();
    }

    /**
     * 批量添加文档
     * 自动分批处理,避免大批量数据导致的内存溢出
     */
    public BulkResult batchAdd(String indexName, List<?> dataList) {
        ParamValidator.validateIndexName(indexName);

        if (CollectionUtils.isEmpty(dataList)) {
            return BulkResult.empty();
        }

        int batchSize = properties.getBulkBatchSize();
        int totalSize = dataList.size();
        int batchCount = (int) Math.ceil((double) totalSize / batchSize);

        BulkResult.Builder resultBuilder = BulkResult.builder();

        for (int i = 0; i < batchCount; i++) {
            int fromIndex = i * batchSize;
            int toIndex = Math.min(fromIndex + batchSize, totalSize);
            List<?> batchData = dataList.subList(fromIndex, toIndex);

            try {
                BulkResult batchResult = doBatchAdd(indexName, batchData);
                resultBuilder.merge(batchResult);

                log.debug("批量添加进度: {}/{}, 成功: {}, 失败: {}",
                        toIndex, totalSize,
                        batchResult.getSuccessCount(),
                        batchResult.getFailureCount());
            } catch (Exception e) {
                log.error("批量添加第 {} 批失败", i + 1, e);
                resultBuilder.addFailure(batchData.size(), e.getMessage());
            }
        }

        BulkResult result = resultBuilder.build();
        log.info("批量添加完成, index: {}, 总数: {}, 成功: {}, 失败: {}",
                indexName, totalSize, result.getSuccessCount(), result.getFailureCount());

        return result;
    }

    private BulkResult doBatchAdd(String indexName, List<?> dataList) throws IOException {
        BulkRequest bulkRequest = new BulkRequest();

        for (Object data : dataList) {
            IndexRequest request = new IndexRequest(indexName);
            request.source(JSON.toJSONString(data), XContentType.JSON);

            if (properties.getEnableType()) {
                request.type(properties.getDefaultType());
            }

            bulkRequest.add(request);
        }

        // 设置刷新策略
        setRefreshPolicy(bulkRequest);

        BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT);

        return parseBulkResponse(response);
    }

    /**
     * 批量添加文档(带ID)
     */
    public BulkResult batchAddWithId(String indexName, Map<String, Object> dataMap) {
        ParamValidator.validateIndexName(indexName);

        if (CollectionUtils.isEmpty(dataMap)) {
            return BulkResult.empty();
        }

        int batchSize = properties.getBulkBatchSize();
        List<Map.Entry<String, Object>> entries = new ArrayList<>(dataMap.entrySet());
        int totalSize = entries.size();
        int batchCount = (int) Math.ceil((double) totalSize / batchSize);

        BulkResult.Builder resultBuilder = BulkResult.builder();

        for (int i = 0; i < batchCount; i++) {
            int fromIndex = i * batchSize;
            int toIndex = Math.min(fromIndex + batchSize, totalSize);
            List<Map.Entry<String, Object>> batchEntries = entries.subList(fromIndex, toIndex);

            try {
                BulkResult batchResult = doBatchAddWithId(indexName, batchEntries);
                resultBuilder.merge(batchResult);
            } catch (Exception e) {
                log.error("批量添加第 {} 批失败", i + 1, e);
                resultBuilder.addFailure(batchEntries.size(), e.getMessage());
            }
        }

        BulkResult result = resultBuilder.build();
        log.info("批量添加完成, index: {}, 总数: {}, 成功: {}, 失败: {}",
                indexName, totalSize, result.getSuccessCount(), result.getFailureCount());

        return result;
    }

    private BulkResult doBatchAddWithId(String indexName,
                                         List<Map.Entry<String, Object>> entries)
            throws IOException {
        BulkRequest bulkRequest = new BulkRequest();

        for (Map.Entry<String, Object> entry : entries) {
            IndexRequest request = new IndexRequest(indexName);
            request.id(entry.getKey());
            request.source(JSON.toJSONString(entry.getValue()), XContentType.JSON);

            if (properties.getEnableType()) {
                request.type(properties.getDefaultType());
            }

            bulkRequest.add(request);
        }

        setRefreshPolicy(bulkRequest);

        BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT);

        return parseBulkResponse(response);
    }

    // ==================== 查询操作 ====================

    /**
     * 根据ID查询
     */
    public <T> Optional<T> getById(String indexName, String id, Class<T> clazz) {
        ParamValidator.validateIndexName(indexName);
        ParamValidator.validateDocumentId(id);

        return RetryUtil.executeWithRetry(
                () -> doGetById(indexName, id, clazz),
                properties.getRetryTimes(),
                properties.getRetryInterval(),
                RetryUtil.defaultRetryCondition(),
                "GET_BY_ID"
        );
    }

    private <T> Optional<T> doGetById(String indexName, String id, Class<T> clazz)
            throws IOException {
        GetRequest request = new GetRequest(indexName, id);
        GetResponse response = client.get(request, RequestOptions.DEFAULT);

        if (!response.isExists()) {
            return Optional.empty();
        }

        String source = response.getSourceAsString();
        if (source == null || source.isEmpty()) {
            return Optional.empty();
        }

        T obj = JSON.parseObject(source, clazz);
        return Optional.of(obj);
    }

    /**
     * 批量根据ID查询
     */
    public <T> Map<String, T> getByIds(String indexName, List<String> ids, Class<T> clazz) {
        ParamValidator.validateIndexName(indexName);

        if (CollectionUtils.isEmpty(ids)) {
            return Collections.emptyMap();
        }

        try {
            MultiGetRequest request = new MultiGetRequest();
            for (String id : ids) {
                request.add(new MultiGetRequest.Item(indexName, id));
            }

            MultiGetResponse response = client.mget(request, RequestOptions.DEFAULT);

            Map<String, T> result = new HashMap<>();
            for (MultiGetItemResponse itemResponse : response.getResponses()) {
                if (!itemResponse.isFailed() && itemResponse.getResponse().isExists()) {
                    String source = itemResponse.getResponse().getSourceAsString();
                    if (source != null && !source.isEmpty()) {
                        T obj = JSON.parseObject(source, clazz);
                        result.put(itemResponse.getId(), obj);
                    }
                }
            }

            return result;
        } catch (IOException e) {
            throw new ElasticsearchException("GET_BY_IDS", indexName, e);
        }
    }

    /**
     * 判断文档是否存在
     */
    public boolean exists(String indexName, String id) {
        ParamValidator.validateIndexName(indexName);
        ParamValidator.validateDocumentId(id);

        try {
            GetRequest request = new GetRequest(indexName, id);
            request.fetchSourceContext(
                    org.elasticsearch.search.fetch.subphase.FetchSourceContext.DO_NOT_FETCH_SOURCE);

            GetResponse response = client.get(request, RequestOptions.DEFAULT);
            return response.isExists();
        } catch (IOException e) {
            throw new ElasticsearchException("EXISTS", indexName, e);
        }
    }

    // ==================== 更新操作 ====================

    /**
     * 更新文档(全量更新,不存在则新增)
     */
    public boolean upsert(String indexName, String id, Object data) {
        ParamValidator.validateIndexName(indexName);
        ParamValidator.validateDocumentId(id);
        ParamValidator.notNull(data, "更新数据不能为空");

        return RetryUtil.executeWithRetry(
                () -> doUpsert(indexName, id, data),
                properties.getRetryTimes(),
                properties.getRetryInterval(),
                RetryUtil.defaultRetryCondition(),
                "UPSERT"
        );
    }

    private boolean doUpsert(String indexName, String id, Object data) throws IOException {
        UpdateRequest request = new UpdateRequest(indexName, id);
        request.doc(JSON.toJSONString(data), XContentType.JSON);
        request.docAsUpsert(true);

        setRefreshPolicy(request);

        UpdateResponse response = client.update(request, RequestOptions.DEFAULT);

        log.info("更新文档成功, index: {}, id: {}, result: {}",
                indexName, id, response.getResult());
        return response.getResult() != DocWriteResponse.Result.NOOP;
    }

    /**
     * 部分更新文档
     */
    public boolean partialUpdate(String indexName, String id, Map<String, Object> fields) {
        ParamValidator.validateIndexName(indexName);
        ParamValidator.validateDocumentId(id);

        if (CollectionUtils.isEmpty(fields)) {
            return false;
        }

        return RetryUtil.executeWithRetry(
                () -> doPartialUpdate(indexName, id, fields),
                properties.getRetryTimes(),
                properties.getRetryInterval(),
                RetryUtil.defaultRetryCondition(),
                "PARTIAL_UPDATE"
        );
    }

    private boolean doPartialUpdate(String indexName, String id,
                                     Map<String, Object> fields) throws IOException {
        UpdateRequest request = new UpdateRequest(indexName, id);
        request.doc(fields);

        setRefreshPolicy(request);

        UpdateResponse response = client.update(request, RequestOptions.DEFAULT);

        log.info("部分更新文档成功, index: {}, id: {}", indexName, id);
        return response.getResult() != DocWriteResponse.Result.NOOP;
    }

    /**
     * 使用脚本更新
     */
    public boolean updateByScript(String indexName, String id, String script,
                                   Map<String, Object> params) {
        ParamValidator.validateIndexName(indexName);
        ParamValidator.validateDocumentId(id);
        ParamValidator.notBlank(script, "脚本不能为空");

        try {
            UpdateRequest request = new UpdateRequest(indexName, id);
            Script painlessScript = new Script(
                    ScriptType.INLINE,
                    "painless",
                    script,
                    params != null ? params : Collections.emptyMap()
            );
            request.script(painlessScript);

            setRefreshPolicy(request);

            UpdateResponse response = client.update(request, RequestOptions.DEFAULT);

            log.info("脚本更新文档成功, index: {}, id: {}", indexName, id);
            return response.getResult() != DocWriteResponse.Result.NOOP;
        } catch (IOException e) {
            throw new ElasticsearchException("UPDATE_BY_SCRIPT", indexName, e);
        }
    }

    /**
     * 根据条件批量更新
     */
    public long updateByQuery(String indexName, QueryBuilder query,
                               String script, Map<String, Object> params) {
        ParamValidator.validateIndexName(indexName);
        ParamValidator.notNull(query, "查询条件不能为空");
        ParamValidator.notBlank(script, "脚本不能为空");

        try {
            UpdateByQueryRequest request = new UpdateByQueryRequest(indexName);
            request.setQuery(query);
            request.setScript(new Script(
                    ScriptType.INLINE,
                    "painless",
                    script,
                    params != null ? params : Collections.emptyMap()
            ));
            request.setRefresh(true);

            BulkByScrollResponse response = client.updateByQuery(
                    request, RequestOptions.DEFAULT);

            log.info("根据条件更新文档成功, index: {}, updated: {}",
                    indexName, response.getUpdated());
            return response.getUpdated();
        } catch (IOException e) {
            throw new ElasticsearchException("UPDATE_BY_QUERY", indexName, e);
        }
    }

    /**
     * 批量更新
     */
    public BulkResult batchUpdate(String indexName, Map<String, Object> dataMap) {
        ParamValidator.validateIndexName(indexName);

        if (CollectionUtils.isEmpty(dataMap)) {
            return BulkResult.empty();
        }

        try {
            BulkRequest bulkRequest = new BulkRequest();

            dataMap.forEach((id, data) -> {
                UpdateRequest request = new UpdateRequest(indexName, id);
                request.doc(JSON.toJSONString(data), XContentType.JSON);
                request.docAsUpsert(true);
                bulkRequest.add(request);
            });

            setRefreshPolicy(bulkRequest);

            BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT);

            BulkResult result = parseBulkResponse(response);
            log.info("批量更新完成, index: {}, 成功: {}, 失败: {}",
                    indexName, result.getSuccessCount(), result.getFailureCount());

            return result;
        } catch (IOException e) {
            throw new ElasticsearchException("BATCH_UPDATE", indexName, e);
        }
    }

    // ==================== 删除操作 ====================

    /**
     * 删除文档
     */
    public boolean delete(String indexName, String id) {
        ParamValidator.validateIndexName(indexName);
        ParamValidator.validateDocumentId(id);

        return RetryUtil.executeWithRetry(
                () -> doDelete(indexName, id),
                properties.getRetryTimes(),
                properties.getRetryInterval(),
                RetryUtil.defaultRetryCondition(),
                "DELETE"
        );
    }

    private boolean doDelete(String indexName, String id) throws IOException {
        DeleteRequest request = new DeleteRequest(indexName, id);

        setRefreshPolicy(request);

        DeleteResponse response = client.delete(request, RequestOptions.DEFAULT);

        log.info("删除文档成功, index: {}, id: {}", indexName, id);
        return response.getResult() == DocWriteResponse.Result.DELETED;
    }

    /**
     * 批量删除
     */
    public BulkResult batchDelete(String indexName, List<String> ids) {
        ParamValidator.validateIndexName(indexName);

        if (CollectionUtils.isEmpty(ids)) {
            return BulkResult.empty();
        }

        try {
            BulkRequest bulkRequest = new BulkRequest();

            for (String id : ids) {
                DeleteRequest request = new DeleteRequest(indexName, id);
                bulkRequest.add(request);
            }

            setRefreshPolicy(bulkRequest);

            BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT);

            BulkResult result = parseBulkResponse(response);
            log.info("批量删除完成, index: {}, 成功: {}, 失败: {}",
                    indexName, result.getSuccessCount(), result.getFailureCount());

            return result;
        } catch (IOException e) {
            throw new ElasticsearchException("BATCH_DELETE", indexName, e);
        }
    }

    /**
     * 根据条件删除
     */
    public long deleteByQuery(String indexName, QueryBuilder query) {
        ParamValidator.validateIndexName(indexName);
        ParamValidator.notNull(query, "查询条件不能为空");

        try {
            DeleteByQueryRequest request = new DeleteByQueryRequest(indexName);
            request.setQuery(query);
            request.setRefresh(true);

            BulkByScrollResponse response = client.deleteByQuery(
                    request, RequestOptions.DEFAULT);

            log.info("根据条件删除文档成功, index: {}, deleted: {}",
                    indexName, response.getDeleted());
            return response.getDeleted();
        } catch (IOException e) {
            throw new ElasticsearchException("DELETE_BY_QUERY", indexName, e);
        }
    }

    // ==================== 内部工具方法 ====================

    /**
     * 设置刷新策略
     */
    private void setRefreshPolicy(IndexRequest request) {
        String policy = properties.getBulkRefreshPolicy();
        if ("immediate".equals(policy)) {
            request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE);
        } else if ("wait_for".equals(policy)) {
            request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL);
        }
    }

    private void setRefreshPolicy(UpdateRequest request) {
        String policy = properties.getBulkRefreshPolicy();
        if ("immediate".equals(policy)) {
            request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE);
        } else if ("wait_for".equals(policy)) {
            request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL);
        }
    }

    private void setRefreshPolicy(DeleteRequest request) {
        String policy = properties.getBulkRefreshPolicy();
        if ("immediate".equals(policy)) {
            request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE);
        } else if ("wait_for".equals(policy)) {
            request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL);
        }
    }

    private void setRefreshPolicy(BulkRequest request) {
        String policy = properties.getBulkRefreshPolicy();
        if ("immediate".equals(policy)) {
            request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE);
        } else if ("wait_for".equals(policy)) {
            request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL);
        }
    }

    /**
     * 解析批量操作响应
     */
    private BulkResult parseBulkResponse(BulkResponse response) {
        BulkResult.Builder builder = BulkResult.builder();

        for (BulkItemResponse itemResponse : response.getItems()) {
            if (itemResponse.isFailed()) {
                builder.addFailure(itemResponse.getId(),
                        itemResponse.getFailureMessage());
            } else {
                builder.addSuccess(itemResponse.getId());
            }
        }

        return builder.build();
    }

    // ==================== 批量操作结果 ====================

    /**
     * 批量操作结果
     */
    @lombok.Data
    public static class BulkResult {
        private int successCount;
        private int failureCount;
        private List<String> successIds;
        private List<FailureItem> failures;

        @lombok.Data
        @lombok.AllArgsConstructor
        public static class FailureItem {
            private String id;
            private String reason;
        }

        public static BulkResult empty() {
            BulkResult result = new BulkResult();
            result.successCount = 0;
            result.failureCount = 0;
            result.successIds = Collections.emptyList();
            result.failures = Collections.emptyList();
            return result;
        }

        public boolean hasFailures() {
            return failureCount > 0;
        }

        public boolean isAllSuccess() {
            return failureCount == 0;
        }

        public static Builder builder() {
            return new Builder();
        }

        public static class Builder {
            private List<String> successIds = new ArrayList<>();
            private List<FailureItem> failures = new ArrayList<>();

            public Builder addSuccess(String id) {
                successIds.add(id);
                return this;
            }

            public Builder addFailure(String id, String reason) {
                failures.add(new FailureItem(id, reason));
                return this;
            }

            public Builder addFailure(int count, String reason) {
                for (int i = 0; i < count; i++) {
                    failures.add(new FailureItem(null, reason));
                }
                return this;
            }

            public Builder merge(BulkResult other) {
                if (other.successIds != null) {
                    successIds.addAll(other.successIds);
                }
                if (other.failures != null) {
                    failures.addAll(other.failures);
                }
                return this;
            }

            public BulkResult build() {
                BulkResult result = new BulkResult();
                result.successIds = successIds;
                result.failures = failures;
                result.successCount = successIds.size();
                result.failureCount = failures.size();
                return result;
            }
        }
    }
}

9.12 滚动查询结果类

package com.example.elasticsearch.result;

import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;

import java.util.List;

/**
 * 滚动查询结果
 */
@Data
@NoArgsConstructor
@AllArgsConstructor
public class ScrollResult<T> {

    /**
     * 滚动ID(用于获取下一批数据)
     */
    private String scrollId;

    /**
     * 当前批次数据
     */
    private List<T> records;

    /**
     * 总记录数
     */
    private Long total;

    /**
     * 是否是最后一批(没有更多数据)
     */
    private Boolean finished;

    /**
     * 判断是否还有更多数据
     */
    public boolean hasMore() {
        return !finished && records != null && !records.isEmpty();
    }
}

9.13 索引操作服务

package com.example.elasticsearch.core;

import com.example.elasticsearch.config.ElasticsearchProperties;
import com.example.elasticsearch.exception.ElasticsearchException;
import lombok.extern.slf4j.Slf4j;
import org.elasticsearch.action.admin.indices.delete.DeleteIndexRequest;
import org.elasticsearch.action.support.master.AcknowledgedResponse;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.client.indices.*;
import org.elasticsearch.common.settings.Settings;
import org.elasticsearch.common.xcontent.XContentType;

import java.io.IOException;
import java.util.Map;

/**
 * Elasticsearch 索引操作服务
 */
@Slf4j
public class ElasticsearchIndexService {

    private final RestHighLevelClient client;
    private final ElasticsearchProperties properties;

    public ElasticsearchIndexService(RestHighLevelClient client,
                                      ElasticsearchProperties properties) {
        this.client = client;
        this.properties = properties;
    }

    /**
     * 创建索引
     */
    public boolean createIndex(String indexName) {
        return createIndex(indexName, null, null);
    }

    /**
     * 创建索引(带 Mapping)
     */
    public boolean createIndex(String indexName, String mapping) {
        try {
            CreateIndexRequest request = new CreateIndexRequest(indexName);

            if (mapping != null) {
                request.source(mapping, XContentType.JSON);
            }

            CreateIndexResponse response = client.indices()
                    .create(request, RequestOptions.DEFAULT);

            log.info("创建索引成功, index: {}", indexName);
            return response.isAcknowledged();
        } catch (IOException e) {
            log.error("创建索引失败", e);
            throw new ElasticsearchException("CREATE_INDEX", indexName, e);
        }
    }

    /**
     * 创建索引(带 Settings 和 Mapping)
     */
    public boolean createIndex(String indexName,
                                Map<String, Object> settings,
                                Map<String, Object> mapping) {
        try {
            CreateIndexRequest request = new CreateIndexRequest(indexName);

            if (settings != null) {
                request.settings(settings);
            }

            if (mapping != null) {
                request.mapping(mapping);
            }

            CreateIndexResponse response = client.indices()
                    .create(request, RequestOptions.DEFAULT);

            log.info("创建索引成功, index: {}", indexName);
            return response.isAcknowledged();
        } catch (IOException e) {
            log.error("创建索引失败", e);
            throw new ElasticsearchException("CREATE_INDEX", indexName, e);
        }
    }

    /**
     * 判断索引是否存在
     */
    public boolean existsIndex(String indexName) {
        try {
            GetIndexRequest request = new GetIndexRequest(indexName);
            return client.indices().exists(request, RequestOptions.DEFAULT);
        } catch (IOException e) {
            log.error("判断索引是否存在失败", e);
            throw new ElasticsearchException("EXISTS_INDEX", indexName, e);
        }
    }

    /**
     * 删除索引
     */
    public boolean deleteIndex(String indexName) {
        try {
            DeleteIndexRequest request = new DeleteIndexRequest(indexName);
            AcknowledgedResponse response = client.indices()
                    .delete(request, RequestOptions.DEFAULT);

            log.info("删除索引成功, index: {}", indexName);
            return response.isAcknowledged();
        } catch (IOException e) {
            log.error("删除索引失败", e);
            throw new ElasticsearchException("DELETE_INDEX", indexName, e);
        }
    }

    /**
     * 更新 Mapping
     */
    public boolean updateMapping(String indexName, Map<String, Object> properties) {
        try {
            PutMappingRequest request = new PutMappingRequest(indexName);
            request.source(Map.of("properties", properties));

            AcknowledgedResponse response = client.indices()
                    .putMapping(request, RequestOptions.DEFAULT);

            log.info("更新 Mapping 成功, index: {}", indexName);
            return response.isAcknowledged();
        } catch (IOException e) {
            log.error("更新 Mapping 失败", e);
            throw new ElasticsearchException("UPDATE_MAPPING", indexName, e);
        }
    }

    /**
     * 获取索引信息
     */
    public GetIndexResponse getIndex(String indexName) {
        try {
            GetIndexRequest request = new GetIndexRequest(indexName);
            return client.indices().get(request, RequestOptions.DEFAULT);
        } catch (IOException e) {
            log.error("获取索引信息失败", e);
            throw new ElasticsearchException("GET_INDEX", indexName, e);
        }
    }

    /**
     * 刷新索引
     */
    public void refreshIndex(String... indexNames) {
        try {
            org.elasticsearch.client.indices.RefreshRequest request =
                    new org.elasticsearch.client.indices.RefreshRequest(indexNames);
            client.indices().refresh(request, RequestOptions.DEFAULT);
            log.info("刷新索引成功, indices: {}", String.join(",", indexNames));
        } catch (IOException e) {
            log.error("刷新索引失败", e);
            throw new ElasticsearchException("REFRESH_INDEX",
                    String.join(",", indexNames), e);
        }
    }

    /**
     * 创建索引别名
     */
    public boolean createAlias(String indexName, String aliasName) {
        try {
            var request = new org.elasticsearch.client.indices.IndicesAliasesRequest();
            var action = new org.elasticsearch.client.indices.IndicesAliasesRequest
                    .AliasActions(org.elasticsearch.client.indices.IndicesAliasesRequest
                    .AliasActions.Type.ADD)
                    .index(indexName)
                    .alias(aliasName);
            request.addAliasAction(action);

            var response = client.indices().updateAliases(request, RequestOptions.DEFAULT);

            log.info("创建索引别名成功, index: {}, alias: {}", indexName, aliasName);
            return response.isAcknowledged();
        } catch (IOException e) {
            log.error("创建索引别名失败", e);
            throw new ElasticsearchException("CREATE_ALIAS", indexName, e);
        }
    }
}

9.14 架构图

7-Blog/后端与微服务/assets/架构图-39a2cd40
  1. 链式调用 - EsQueryBuilder 支持流畅的链式构建
  2. 类型安全 - 泛型支持,自动转换结果类型
  3. 功能完整 - 覆盖索引、文档、查询、聚合操作
  4. 高亮支持 - 内置高亮配置和结果解析
  5. 分页封装 - 统一的分页结果对象
  6. 连接池 - 支持集群、认证、连接池配置
  7. 版本兼容 - 支持 ES 6.x/7.x type 兼容
  8. 异常处理 - 统一的异常封装

十、完整使用文档

快速开始

添加依赖

<!-- pom.xml -->
<dependencies>
    <!-- ES Starter -->
    <dependency>
        <groupId>com.example</groupId>
        <artifactId>elasticsearch-spring-boot-starter</artifactId>
        <version>1.0.0</version>
    </dependency>

    <!-- 或者使用官方依赖 -->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-data-elasticsearch</artifactId>
    </dependency>

    <!-- FastJSON -->
    <dependency>
        <groupId>com.alibaba</groupId>
        <artifactId>fastjson</artifactId>
        <version>1.2.83</version>
    </dependency>
</dependencies>

配置文件

# application.yml
elasticsearch:
  enabled: true
  nodes:
    - localhost:9200
  username: elastic          # 可选
  password: your_password    # 可选
  scheme: http

  # 连接配置
  connect-timeout: 5000
  socket-timeout: 30000
  max-conn-total: 100
  max-conn-per-route: 50

  # 重试配置
  retry-times: 3
  retry-interval: 1000

  # 批量操作配置
  bulk-batch-size: 1000
  bulk-refresh-policy: none  # none, immediate, wait_for

  # 日志配置
  print-dsl: true            # 开发环境建议开启
  slow-query-threshold: 3000 # 慢查询阈值(ms)

  # 版本兼容
  enable-type: false         # ES 7.x 设为 false

创建实体类

package com.example.demo.entity;

import lombok.Data;
import java.util.Date;

/**
 * 节目实体
 */
@Data
public class Program {

    private Long id;

    /**
     * 节目标题
     */
    private String title;

    /**
     * 演员
     */
    private String actor;

    /**
     * 分类ID
     */
    private Long categoryId;

    /**
     * 分类名称
     */
    private String categoryName;

    /**
     * 地区ID
     */
    private Long areaId;

    /**
     * 地区名称
     */
    private String areaName;

    /**
     * 演出时间
     */
    private Date showTime;

    /**
     * 最低价格
     */
    private Double minPrice;

    /**
     * 最高价格
     */
    private Double maxPrice;

    /**
     * 状态:1-上架, 0-下架
     */
    private Integer status;

    /**
     * 封面图片
     */
    private String coverImage;

    /**
     * 创建时间
     */
    private Date createTime;

    /**
     * 更新时间
     */
    private Date updateTime;
}

创建搜索参数类

package com.example.demo.dto;

import lombok.Data;
import java.util.Date;

/**
 * 节目搜索参数
 */
@Data
public class ProgramSearchParam {

    /**
     * 关键词(搜索标题和演员)
     */
    private String keyword;

    /**
     * 分类ID
     */
    private Long categoryId;

    /**
     * 地区ID
     */
    private Long areaId;

    /**
     * 状态
     */
    private Integer status;

    /**
     * 最低价格
     */
    private Double minPrice;

    /**
     * 最高价格
     */
    private Double maxPrice;

    /**
     * 开始时间
     */
    private Date startTime;

    /**
     * 结束时间
     */
    private Date endTime;

    /**
     * 页码
     */
    private Integer pageNum = 1;

    /**
     * 每页大小
     */
    private Integer pageSize = 10;

    /**
     * 排序字段
     */
    private String sortField = "showTime";

    /**
     * 排序方式:asc, desc
     */
    private String sortOrder = "asc";
}

完整业务示例

package com.example.demo.service;

import com.example.demo.dto.ProgramSearchParam;
import com.example.demo.entity.Program;
import com.example.elasticsearch.core.ElasticsearchDocumentService;
import com.example.elasticsearch.core.ElasticsearchDocumentService.BulkResult;
import com.example.elasticsearch.core.ElasticsearchIndexService;
import com.example.elasticsearch.core.ElasticsearchService;
import com.example.elasticsearch.query.EsQueryBuilder;
import com.example.elasticsearch.result.AggregationResult;
import com.example.elasticsearch.result.PageResult;
import com.example.elasticsearch.result.ScrollResult;
import com.example.elasticsearch.result.SearchResult;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.elasticsearch.search.sort.SortOrder;
import org.springframework.stereotype.Service;

import javax.annotation.PostConstruct;
import java.util.*;
import java.util.function.Consumer;
import java.util.stream.Collectors;

/**
 * 节目搜索服务 - 完整示例
 */
@Slf4j
@Service
@RequiredArgsConstructor
public class ProgramSearchService {

    private final ElasticsearchService esService;
    private final ElasticsearchDocumentService documentService;
    private final ElasticsearchIndexService indexService;

    private static final String INDEX_NAME = "program";

    // ==================== 索引管理 ====================

    /**
     * 初始化索引(应用启动时调用)
     */
    @PostConstruct
    public void initIndex() {
        if (!indexService.existsIndex(INDEX_NAME)) {
            String mapping = """
                {
                    "mappings": {
                        "properties": {
                            "id": { "type": "long" },
                            "title": {
                                "type": "text",
                                "analyzer": "ik_max_word",
                                "search_analyzer": "ik_smart",
                                "fields": {
                                    "keyword": { "type": "keyword" }
                                }
                            },
                            "actor": {
                                "type": "text",
                                "analyzer": "ik_max_word",
                                "fields": {
                                    "keyword": { "type": "keyword" }
                                }
                            },
                            "categoryId": { "type": "long" },
                            "categoryName": { "type": "keyword" },
                            "areaId": { "type": "long" },
                            "areaName": { "type": "keyword" },
                            "showTime": { "type": "date" },
                            "minPrice": { "type": "double" },
                            "maxPrice": { "type": "double" },
                            "status": { "type": "integer" },
                            "coverImage": { "type": "keyword", "index": false },
                            "createTime": { "type": "date" },
                            "updateTime": { "type": "date" }
                        }
                    },
                    "settings": {
                        "number_of_shards": 3,
                        "number_of_replicas": 1,
                        "refresh_interval": "1s"
                    }
                }
                """;
            indexService.createIndex(INDEX_NAME, mapping);
            log.info("索引 [{}] 创建成功", INDEX_NAME);
        }
    }

    /**
     * 重建索引
     */
    public void rebuildIndex() {
        // 删除旧索引
        if (indexService.existsIndex(INDEX_NAME)) {
            indexService.deleteIndex(INDEX_NAME);
        }
        // 重新初始化
        initIndex();
    }

    // ==================== 文档操作 ====================

    /**
     * 添加/更新节目
     */
    public void saveProgram(Program program) {
        program.setUpdateTime(new Date());
        if (program.getCreateTime() == null) {
            program.setCreateTime(new Date());
        }
        documentService.upsert(INDEX_NAME, String.valueOf(program.getId()), program);
    }

    /**
     * 批量添加节目
     */
    public BulkResult batchSavePrograms(List<Program> programs) {
        Date now = new Date();
        programs.forEach(p -> {
            p.setUpdateTime(now);
            if (p.getCreateTime() == null) {
                p.setCreateTime(now);
            }
        });
        return documentService.batchAdd(INDEX_NAME, programs);
    }

    /**
     * 批量添加节目(带ID)
     */
    public BulkResult batchSaveProgramsWithId(List<Program> programs) {
        Date now = new Date();
        Map<String, Object> dataMap = new HashMap<>();

        programs.forEach(p -> {
            p.setUpdateTime(now);
            if (p.getCreateTime() == null) {
                p.setCreateTime(now);
            }
            dataMap.put(String.valueOf(p.getId()), p);
        });

        return documentService.batchAddWithId(INDEX_NAME, dataMap);
    }

    /**
     * 根据ID查询节目
     */
    public Optional<Program> getProgramById(Long id) {
        return documentService.getById(INDEX_NAME, String.valueOf(id), Program.class);
    }

    /**
     * 批量查询节目
     */
    public Map<Long, Program> getProgramsByIds(List<Long> ids) {
        List<String> strIds = ids.stream()
                .map(String::valueOf)
                .collect(Collectors.toList());

        Map<String, Program> resultMap = documentService.getByIds(
                INDEX_NAME, strIds, Program.class);

        // 转换 key 类型
        return resultMap.entrySet().stream()
                .collect(Collectors.toMap(
                        e -> Long.parseLong(e.getKey()),
                        Map.Entry::getValue
                ));
    }

    /**
     * 删除节目
     */
    public void deleteProgram(Long id) {
        documentService.delete(INDEX_NAME, String.valueOf(id));
    }

    /**
     * 批量删除节目
     */
    public BulkResult batchDeletePrograms(List<Long> ids) {
        List<String> strIds = ids.stream()
                .map(String::valueOf)
                .collect(Collectors.toList());
        return documentService.batchDelete(INDEX_NAME, strIds);
    }

    /**
     * 更新节目状态
     */
    public boolean updateProgramStatus(Long id, Integer status) {
        Map<String, Object> fields = new HashMap<>();
        fields.put("status", status);
        fields.put("updateTime", new Date());
        return documentService.partialUpdate(INDEX_NAME, String.valueOf(id), fields);
    }

    /**
     * 批量更新节目状态(使用脚本)
     */
    public long batchUpdateStatus(List<Long> ids, Integer status) {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                .filterTerms("id", ids);

        String script = "ctx._source.status = params.status; " +
                        "ctx._source.updateTime = params.updateTime";

        Map<String, Object> params = new HashMap<>();
        params.put("status", status);
        params.put("updateTime", System.currentTimeMillis());

        return documentService.updateByQuery(
                INDEX_NAME, queryBuilder.getBoolQuery(), script, params);
    }

    // ==================== 基础搜索 ====================

    /**
     * 根据分类查询节目列表
     */
    public List<Program> findByCategory(Long categoryId) {
        return esService.query(INDEX_NAME, "categoryId", categoryId, Program.class);
    }

    /**
     * 根据多条件查询
     */
    public List<Program> findByConditions(Long categoryId, Long areaId, Integer status) {
        Map<String, Object> params = new HashMap<>();
        params.put("categoryId", categoryId);
        params.put("areaId", areaId);
        params.put("status", status);

        return esService.query(INDEX_NAME, params, Program.class);
    }

    /**
     * 简单分页查询
     */
    public PageResult<Program> findPage(int pageNum, int pageSize) {
        return esService.queryPage(INDEX_NAME, null, pageNum, pageSize, Program.class);
    }

    // ==================== 高级搜索 ====================

    /**
     * 关键词搜索(标题 + 演员)
     */
    public List<Program> searchByKeyword(String keyword) {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                .shouldMatch(keyword, "title", "actor")
                .filterTerm("status", 1) // 只搜索上架的
                .sort("showTime", SortOrder.ASC);

        return esService.search(queryBuilder, Program.class);
    }

    /**
     * 关键词搜索(带高亮)
     */
    public List<SearchResult<Program>> searchWithHighlight(String keyword) {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                .shouldMatch(keyword, "title", "actor")
                .filterTerm("status", 1)
                .highlight("<span class='highlight'>", "</span>", "title", "actor")
                .sort("showTime", SortOrder.ASC);

        return esService.searchWithHighlight(queryBuilder, Program.class);
    }

    /**
     * 复杂条件搜索 + 分页
     */
    public PageResult<Program> search(ProgramSearchParam param) {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                // 关键词搜索(标题或演员)
                .shouldMatch(param.getKeyword(), "title", "actor")
                // 分类过滤
                .filterTerm("categoryId", param.getCategoryId())
                // 地区过滤
                .filterTerm("areaId", param.getAreaId())
                // 状态过滤
                .filterTerm("status", param.getStatus())
                // 价格范围
                .filterRange("minPrice", param.getMinPrice(), param.getMaxPrice())
                // 时间范围
                .filterRange("showTime",
                        param.getStartTime() != null ? param.getStartTime().getTime() : null,
                        param.getEndTime() != null ? param.getEndTime().getTime() : null)
                // 排序
                .sort(param.getSortField(),
                        "desc".equalsIgnoreCase(param.getSortOrder())
                                ? SortOrder.DESC : SortOrder.ASC);

        return esService.searchPage(queryBuilder,
                param.getPageNum(), param.getPageSize(), Program.class);
    }

    /**
     * 复杂条件搜索 + 分页 + 高亮
     */
    public PageResult<SearchResult<Program>> searchWithHighlight(ProgramSearchParam param) {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                .shouldMatch(param.getKeyword(), "title", "actor")
                .filterTerm("categoryId", param.getCategoryId())
                .filterTerm("areaId", param.getAreaId())
                .filterTerm("status", param.getStatus())
                .filterRange("minPrice", param.getMinPrice(), param.getMaxPrice())
                .filterRange("showTime",
                        param.getStartTime() != null ? param.getStartTime().getTime() : null,
                        param.getEndTime() != null ? param.getEndTime().getTime() : null)
                .highlight("title", "actor")
                .sort(param.getSortField(),
                        "desc".equalsIgnoreCase(param.getSortOrder())
                                ? SortOrder.DESC : SortOrder.ASC);

        return esService.searchPageWithHighlight(queryBuilder,
                param.getPageNum(), param.getPageSize(), Program.class);
    }

    // ==================== 深度分页(Search After) ====================

    /**
     * 使用 Search After 进行深度分页
     * 适用于页码超过 1000 的场景
     *
     * @param param            搜索参数
     * @param searchAfterValue 上一页最后一条记录的排序值
     * @return 搜索结果(包含 sortValues 用于下一页查询)
     */
    public List<SearchResult<Program>> searchAfter(ProgramSearchParam param,
                                                     Object[] searchAfterValue) {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                .shouldMatch(param.getKeyword(), "title", "actor")
                .filterTerm("categoryId", param.getCategoryId())
                .filterTerm("status", param.getStatus())
                // 必须有排序字段,且最后一个字段应该是唯一的(如 _id)
                .sort(param.getSortField(),
                        "desc".equalsIgnoreCase(param.getSortOrder())
                                ? SortOrder.DESC : SortOrder.ASC)
                .sort("id", SortOrder.ASC); // 使用 id 保证唯一性

        return esService.searchAfter(queryBuilder, searchAfterValue,
                param.getPageSize(), Program.class);
    }

    // ==================== 滚动查询(大数据导出) ====================

    /**
     * 滚动查询导出所有数据
     * 适用于数据导出场景
     *
     * @param param    搜索参数
     * @param consumer 数据消费者
     */
    public void scrollExport(ProgramSearchParam param, Consumer<List<Program>> consumer) {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                .shouldMatch(param.getKeyword(), "title", "actor")
                .filterTerm("categoryId", param.getCategoryId())
                .filterTerm("status", param.getStatus())
                .sort("id", SortOrder.ASC); // 滚动查询必须有排序

        int batchSize = 1000;
        int scrollTime = 5; // 滚动上下文保持时间(分钟)

        // 初始化滚动查询
        ScrollResult<Program> scrollResult = esService.scrollSearch(
                queryBuilder, scrollTime, batchSize, Program.class);

        try {
            // 处理首批数据
            if (!scrollResult.getRecords().isEmpty()) {
                consumer.accept(scrollResult.getRecords());
            }

            // 继续获取后续数据
            while (scrollResult.hasMore()) {
                scrollResult = esService.scrollNext(
                        scrollResult.getScrollId(), scrollTime, batchSize, Program.class);

                if (!scrollResult.getRecords().isEmpty()) {
                    consumer.accept(scrollResult.getRecords());
                }
            }

            log.info("滚动导出完成, 总数: {}", scrollResult.getTotal());
        } finally {
            // 清除滚动上下文
            esService.clearScroll(scrollResult.getScrollId());
        }
    }

    // ==================== 聚合统计 ====================

    /**
     * 按分类统计节目数量
     */
    public Map<String, Long> countByCategory() {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                .filterTerm("status", 1) // 只统计上架的
                .termsAggregation("category_count", "categoryId", 100);

        Map<String, AggregationResult> aggResult = esService.searchAggregation(queryBuilder);

        AggregationResult categoryAgg = aggResult.get("category_count");
        if (categoryAgg == null || categoryAgg.getBuckets() == null) {
            return Collections.emptyMap();
        }

        return categoryAgg.getBuckets().stream()
                .collect(Collectors.toMap(
                        AggregationResult.BucketData::getKey,
                        AggregationResult.BucketData::getDocCount
                ));
    }

    /**
     * 统计价格区间
     */
    public Map<String, Double> getPriceStats() {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                .filterTerm("status", 1)
                .minAggregation("min_price", "minPrice")
                .maxAggregation("max_price", "maxPrice")
                .avgAggregation("avg_price", "minPrice");

        Map<String, AggregationResult> aggResult = esService.searchAggregation(queryBuilder);

        Map<String, Double> stats = new HashMap<>();

        if (aggResult.containsKey("min_price")) {
            stats.put("minPrice", aggResult.get("min_price").getValue());
        }
        if (aggResult.containsKey("max_price")) {
            stats.put("maxPrice", aggResult.get("max_price").getValue());
        }
        if (aggResult.containsKey("avg_price")) {
            stats.put("avgPrice", aggResult.get("avg_price").getValue());
        }

        return stats;
    }

    /**
     * 综合统计
     */
    public Map<String, Object> getStatistics() {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                .filterTerm("status", 1)
                .termsAggregation("by_category", "categoryId", 50)
                .termsAggregation("by_area", "areaId", 50)
                .avgAggregation("avg_price", "minPrice")
                .minAggregation("min_price", "minPrice")
                .maxAggregation("max_price", "maxPrice");

        Map<String, AggregationResult> aggResult = esService.searchAggregation(queryBuilder);

        Map<String, Object> stats = new HashMap<>();
        stats.put("aggregations", aggResult);
        stats.put("total", esService.count(
                EsQueryBuilder.builder(INDEX_NAME).filterTerm("status", 1)));

        return stats;
    }

    // ==================== 工具方法 ====================

    /**
     * 统计符合条件的数量
     */
    public long count(ProgramSearchParam param) {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                .shouldMatch(param.getKeyword(), "title", "actor")
                .filterTerm("categoryId", param.getCategoryId())
                .filterTerm("areaId", param.getAreaId())
                .filterTerm("status", param.getStatus());

        return esService.count(queryBuilder);
    }

    /**
     * 判断是否存在符合条件的数据
     */
    public boolean exists(ProgramSearchParam param) {
        return count(param) > 0;
    }
}

Controller 示例

package com.example.demo.controller;

import com.example.demo.dto.ProgramSearchParam;
import com.example.demo.entity.Program;
import com.example.demo.service.ProgramSearchService;
import com.example.elasticsearch.core.ElasticsearchDocumentService.BulkResult;
import com.example.elasticsearch.result.PageResult;
import com.example.elasticsearch.result.SearchResult;
import lombok.RequiredArgsConstructor;
import org.springframework.web.bind.annotation.*;

import java.util.List;
import java.util.Map;
import java.util.Optional;

/**
 * 节目搜索 API
 */
@RestController
@RequestMapping("/api/programs")
@RequiredArgsConstructor
public class ProgramController {

    private final ProgramSearchService searchService;

    // ==================== 基础 CRUD ====================

    /**
     * 添加/更新节目
     */
    @PostMapping
    public String saveProgram(@RequestBody Program program) {
        searchService.saveProgram(program);
        return "success";
    }

    /**
     * 批量添加节目
     */
    @PostMapping("/batch")
    public BulkResult batchSavePrograms(@RequestBody List<Program> programs) {
        return searchService.batchSavePrograms(programs);
    }

    /**
     * 根据ID查询
     */
    @GetMapping("/{id}")
    public Optional<Program> getProgram(@PathVariable Long id) {
        return searchService.getProgramById(id);
    }

    /**
     * 批量查询
     */
    @PostMapping("/batch-get")
    public Map<Long, Program> batchGetPrograms(@RequestBody List<Long> ids) {
        return searchService.getProgramsByIds(ids);
    }

    /**
     * 删除节目
     */
    @DeleteMapping("/{id}")
    public String deleteProgram(@PathVariable Long id) {
        searchService.deleteProgram(id);
        return "success";
    }

    /**
     * 更新状态
     */
    @PutMapping("/{id}/status")
    public boolean updateStatus(@PathVariable Long id, @RequestParam Integer status) {
        return searchService.updateProgramStatus(id, status);
    }

    // ==================== 搜索 ====================

    /**
     * 关键词搜索
     */
    @GetMapping("/search")
    public List<Program> search(@RequestParam String keyword) {
        return searchService.searchByKeyword(keyword);
    }

    /**
     * 关键词搜索(带高亮)
     */
    @GetMapping("/search/highlight")
    public List<SearchResult<Program>> searchWithHighlight(@RequestParam String keyword) {
        return searchService.searchWithHighlight(keyword);
    }

    /**
     * 高级搜索(分页)
     */
    @PostMapping("/search/page")
    public PageResult<Program> searchPage(@RequestBody ProgramSearchParam param) {
        return searchService.search(param);
    }

    /**
     * 高级搜索(分页 + 高亮)
     */
    @PostMapping("/search/page-highlight")
    public PageResult<SearchResult<Program>> searchPageWithHighlight(
            @RequestBody ProgramSearchParam param) {
        return searchService.searchWithHighlight(param);
    }

    /**
     * 深度分页(Search After)
     */
    @PostMapping("/search/after")
    public List<SearchResult<Program>> searchAfter(
            @RequestBody ProgramSearchParam param,
            @RequestParam(required = false) Long afterId,
            @RequestParam(required = false) Long afterShowTime) {

        Object[] searchAfter = null;
        if (afterShowTime != null && afterId != null) {
            searchAfter = new Object[]{afterShowTime, afterId};
        }

        return searchService.searchAfter(param, searchAfter);
    }

    // ==================== 统计 ====================

    /**
     * 按分类统计
     */
    @GetMapping("/stats/category")
    public Map<String, Long> countByCategory() {
        return searchService.countByCategory();
    }

    /**
     * 价格统计
     */
    @GetMapping("/stats/price")
    public Map<String, Double> getPriceStats() {
        return searchService.getPriceStats();
    }

    /**
     * 综合统计
     */
    @GetMapping("/stats")
    public Map<String, Object> getStatistics() {
        return searchService.getStatistics();
    }

    // ==================== 索引管理 ====================

    /**
     * 重建索引
     */
    @PostMapping("/index/rebuild")
    public String rebuildIndex() {
        searchService.rebuildIndex();
        return "success";
    }
}

最佳实践指南

索引设计

/**
 * 索引设计最佳实践
 */
public class IndexDesignBestPractice {

    /**
     * 1. 合理设置分片数
     * - 单个分片大小建议 10-50GB
     * - 分片数 = 数据量 / 单分片大小
     */
    public static final int SHARD_COUNT = 3;

    /**
     * 2. 副本数设置
     * - 生产环境至少 1 个副本
     * - 写入密集型可先设为 0,完成后再设为 1
     */
    public static final int REPLICA_COUNT = 1;

    /**
     * 3. Mapping 设计原则
     * - 明确字段类型,避免动态映射
     * - text 用于全文搜索,keyword 用于精确匹配/聚合
     * - 不需要搜索的字段设置 index: false
     * - 大文本字段考虑使用 store: true
     */
    public static String createOptimizedMapping() {
        return """
            {
                "mappings": {
                    "properties": {
                        "id": { "type": "long" },
                        "title": {
                            "type": "text",
                            "analyzer": "ik_max_word",
                            "search_analyzer": "ik_smart",
                            "fields": {
                                "keyword": {
                                    "type": "keyword",
                                    "ignore_above": 256
                                }
                            }
                        },
                        "description": {
                            "type": "text",
                            "analyzer": "ik_max_word",
                            "index": true,
                            "store": true
                        },
                        "status": { "type": "keyword" },
                        "price": { "type": "scaled_float", "scaling_factor": 100 },
                        "createTime": { "type": "date", "format": "epoch_millis" },
                        "tags": { "type": "keyword" },
                        "coverImage": { "type": "keyword", "index": false }
                    }
                },
                "settings": {
                    "number_of_shards": 3,
                    "number_of_replicas": 1,
                    "refresh_interval": "1s",
                    "max_result_window": 10000,
                    "analysis": {
                        "analyzer": {
                            "ik_smart_pinyin": {
                                "type": "custom",
                                "tokenizer": "ik_smart",
                                "filter": ["lowercase", "pinyin_filter"]
                            }
                        }
                    }
                }
            }
            """;
    }
}

查询优化

/**
 * 查询优化最佳实践
 */
@Slf4j
@Service
public class QueryOptimizationService {

    @Autowired
    private ElasticsearchService esService;

    private static final String INDEX_NAME = "program";

    /**
     * 1. 使用 filter 而非 must 进行过滤
     * filter 不计算评分,性能更好,且结果可缓存
     */
    public List<Program> optimizedFilterQuery(Long categoryId, Integer status) {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                // ✅ 正确:使用 filter 进行精确匹配
                .filterTerm("categoryId", categoryId)
                .filterTerm("status", status);

        // ❌ 错误:使用 must 进行精确匹配会计算评分,浪费性能
        // .term("categoryId", categoryId)
        // .term("status", status)

        return esService.search(queryBuilder, Program.class);
    }

    /**
     * 2. 合理使用分页,避免深度分页
     * from + size 超过 10000 会报错
     */
    public PageResult<Program> safePageQuery(int pageNum, int pageSize) {
        // 计算深度
        int depth = (pageNum - 1) * pageSize + pageSize;

        if (depth > 10000) {
            // 深度分页使用 search_after
            log.warn("分页深度超过10000,建议使用 searchAfter");
            throw new IllegalArgumentException("分页深度超出限制,请使用 searchAfter");
        }

        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                .page(pageNum, pageSize)
                .sort("id", SortOrder.ASC);

        return esService.searchPage(queryBuilder, pageNum, pageSize, Program.class);
    }

    /**
     * 3. 只返回需要的字段,减少网络传输
     */
    public List<Program> selectiveFieldsQuery(String keyword) {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                .match("title", keyword)
                // 只返回需要的字段
                .includes("id", "title", "minPrice", "showTime")
                // 排除大字段
                .excludes("description", "coverImage");

        return esService.search(queryBuilder, Program.class);
    }

    /**
     * 4. 使用 bool 查询组合多个条件
     */
    public List<Program> complexBoolQuery(String keyword, List<Long> categoryIds,
                                           Double minPrice, Double maxPrice,
                                           List<Integer> excludeStatus) {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                // must: 必须匹配(影响评分)
                .match("title", keyword)
                // filter: 必须匹配(不影响评分,可缓存)
                .filterTerms("categoryId", categoryIds)
                .filterRange("minPrice", minPrice, maxPrice)
                // must_not: 必须不匹配
                .mustNotTerms("status", excludeStatus)
                // should: 可选匹配(有则加分)
                .should(QueryBuilders.matchQuery("actor", keyword))
                .minimumShouldMatch(0);

        return esService.search(queryBuilder, Program.class);
    }

    /**
     * 5. 使用 constant_score 包装不需要评分的查询
     */
    public List<Program> constantScoreQuery(Long categoryId) {
        // 使用原生 QueryBuilder 实现 constant_score
        QueryBuilder query = QueryBuilders.constantScoreQuery(
                QueryBuilders.termQuery("categoryId", categoryId)
        ).boost(1.0f);

        return esService.searchByQueryBuilder(INDEX_NAME, query, 100, Program.class);
    }

    /**
     * 6. 批量查询优化:使用 multi_search
     */
    public Map<String, List<Program>> multiSearch(List<Long> categoryIds) {
        // 注意:需要扩展 ElasticsearchService 支持 multi_search
        // 这里展示思路
        Map<String, List<Program>> result = new HashMap<>();

        for (Long categoryId : categoryIds) {
            EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                    .filterTerm("categoryId", categoryId)
                    .fromSize(0, 10);

            List<Program> programs = esService.search(queryBuilder, Program.class);
            result.put(String.valueOf(categoryId), programs);
        }

        return result;
    }

    /**
     * 7. 高亮优化:限制高亮片段
     */
    public List<SearchResult<Program>> optimizedHighlightQuery(String keyword) {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
                .match("title", keyword)
                .match("description", keyword);

        // 自定义高亮配置
        HighlightBuilder highlightBuilder = new HighlightBuilder()
                .field(new HighlightBuilder.Field("title")
                        .fragmentSize(100)      // 片段大小
                        .numOfFragments(1))     // 片段数量
                .field(new HighlightBuilder.Field("description")
                        .fragmentSize(150)
                        .numOfFragments(2))
                .preTags("<em>")
                .postTags("</em>");

        queryBuilder.setHighlightBuilder(highlightBuilder);

        return esService.searchWithHighlight(queryBuilder, Program.class);
    }
}

写入优化

/**
 * 写入优化最佳实践
 */
@Slf4j
@Service
public class WriteOptimizationService {

    @Autowired
    private ElasticsearchDocumentService documentService;

    @Autowired
    private ElasticsearchIndexService indexService;

    private static final String INDEX_NAME = "program";

    /**
     * 1. 批量写入优化
     * - 合理设置批次大小(通常 1000-5000)
     * - 控制并发写入线程数
     */
    public void optimizedBatchWrite(List<Program> programs) {
        int batchSize = 1000;
        int totalSize = programs.size();

        log.info("开始批量写入,总数: {}", totalSize);
        long startTime = System.currentTimeMillis();

        // 分批处理
        for (int i = 0; i < totalSize; i += batchSize) {
            int endIndex = Math.min(i + batchSize, totalSize);
            List<Program> batch = programs.subList(i, endIndex);

            BulkResult result = documentService.batchAdd(INDEX_NAME, batch);

            if (result.hasFailures()) {
                log.warn("批次 {}-{} 部分失败,失败数: {}",
                        i, endIndex, result.getFailureCount());
            }

            log.info("写入进度: {}/{}", endIndex, totalSize);
        }

        long elapsed = System.currentTimeMillis() - startTime;
        log.info("批量写入完成,耗时: {}ms,速率: {} docs/s",
                elapsed, totalSize * 1000L / elapsed);
    }

    /**
     * 2. 大批量数据导入优化
     * - 临时关闭副本
     * - 增大 refresh_interval
     * - 导入完成后恢复
     */
    public void bulkImportWithOptimization(List<Program> programs) {
        try {
            // 优化索引设置
            optimizeIndexForBulkImport();

            // 执行批量导入
            optimizedBatchWrite(programs);

        } finally {
            // 恢复索引设置
            restoreIndexSettings();
        }
    }

    private void optimizeIndexForBulkImport() {
        Map<String, Object> settings = new HashMap<>();
        settings.put("index.refresh_interval", "-1");      // 关闭自动刷新
        settings.put("index.number_of_replicas", 0);       // 临时关闭副本

        indexService.updateSettings(INDEX_NAME, settings);
        log.info("索引设置已优化");
    }

    private void restoreIndexSettings() {
        Map<String, Object> settings = new HashMap<>();
        settings.put("index.refresh_interval", "1s");      // 恢复刷新间隔
        settings.put("index.number_of_replicas", 1);       // 恢复副本

        indexService.updateSettings(INDEX_NAME, settings);

        // 强制刷新
        indexService.refreshIndex(INDEX_NAME);

        // 强制合并段
        indexService.forceMerge(INDEX_NAME, 1);

        log.info("索引设置已恢复");
    }

    /**
     * 3. 使用异步写入
     */
    @Async
    public CompletableFuture<BulkResult> asyncBatchWrite(List<Program> programs) {
        BulkResult result = documentService.batchAdd(INDEX_NAME, programs);
        return CompletableFuture.completedFuture(result);
    }

    /**
     * 4. 并行批量写入
     */
    public void parallelBatchWrite(List<Program> programs, int threadCount) {
        int batchSize = programs.size() / threadCount;

        ExecutorService executor = Executors.newFixedThreadPool(threadCount);
        List<Future<BulkResult>> futures = new ArrayList<>();

        try {
            for (int i = 0; i < threadCount; i++) {
                int start = i * batchSize;
                int end = (i == threadCount - 1) ? programs.size() : start + batchSize;
                List<Program> batch = programs.subList(start, end);

                Future<BulkResult> future = executor.submit(() ->
                        documentService.batchAdd(INDEX_NAME, batch));
                futures.add(future);
            }

            // 等待所有任务完成
            int totalSuccess = 0;
            int totalFailure = 0;

            for (Future<BulkResult> future : futures) {
                BulkResult result = future.get();
                totalSuccess += result.getSuccessCount();
                totalFailure += result.getFailureCount();
            }

            log.info("并行写入完成,成功: {}, 失败: {}", totalSuccess, totalFailure);

        } catch (Exception e) {
            log.error("并行写入失败", e);
            throw new RuntimeException(e);
        } finally {
            executor.shutdown();
        }
    }
}

缓存策略

/**
 * ES 查询缓存策略
 */
@Slf4j
@Service
public class EsCacheService {

    @Autowired
    private ElasticsearchService esService;

    @Autowired
    private RedisTemplate<String, Object> redisTemplate;

    private static final String CACHE_PREFIX = "es:program:";
    private static final Duration CACHE_TTL = Duration.ofMinutes(5);

    /**
     * 1. 使用 Redis 缓存查询结果
     */
    @SuppressWarnings("unchecked")
    public List<Program> searchWithCache(String keyword) {
        String cacheKey = CACHE_PREFIX + "search:" + DigestUtils.md5Hex(keyword);

        // 先查缓存
        Object cached = redisTemplate.opsForValue().get(cacheKey);
        if (cached != null) {
            log.debug("命中缓存: {}", cacheKey);
            return (List<Program>) cached;
        }

        // 查询 ES
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder("program")
                .match("title", keyword);

        List<Program> result = esService.search(queryBuilder, Program.class);

        // 写入缓存
        if (!result.isEmpty()) {
            redisTemplate.opsForValue().set(cacheKey, result, CACHE_TTL);
        }

        return result;
    }

    /**
     * 2. 热门搜索词预热缓存
     */
    @Scheduled(fixedRate = 300000) // 每5分钟
    public void warmUpCache() {
        List<String> hotKeywords = getHotKeywords();

        for (String keyword : hotKeywords) {
            try {
                searchWithCache(keyword);
            } catch (Exception e) {
                log.warn("预热缓存失败: {}", keyword, e);
            }
        }

        log.info("缓存预热完成,预热 {} 个关键词", hotKeywords.size());
    }

    /**
     * 3. 数据更新时清除相关缓存
     */
    public void invalidateCache(Long programId) {
        // 清除该节目相关的所有缓存
        Set<String> keys = redisTemplate.keys(CACHE_PREFIX + "*");
        if (keys != null && !keys.isEmpty()) {
            redisTemplate.delete(keys);
            log.info("清除 {} 个缓存", keys.size());
        }
    }

    /**
     * 4. 分类统计结果缓存(变化不频繁的数据)
     */
    @Cacheable(value = "es:stats", key = "'category'", unless = "#result == null")
    public Map<String, Long> getCategoryStats() {
        EsQueryBuilder queryBuilder = EsQueryBuilder.builder("program")
                .termsAggregation("category_count", "categoryId", 100);

        Map<String, AggregationResult> aggResult =
                esService.searchAggregation(queryBuilder);

        // 解析结果
        Map<String, Long> stats = new HashMap<>();
        AggregationResult categoryAgg = aggResult.get("category_count");
        if (categoryAgg != null && categoryAgg.getBuckets() != null) {
            for (AggregationResult.BucketData bucket : categoryAgg.getBuckets()) {
                stats.put(bucket.getKey(), bucket.getDocCount());
            }
        }

        return stats;
    }

    private List<String> getHotKeywords() {
        // 从统计服务获取热门搜索词
        return Arrays.asList("演唱会", "话剧", "音乐剧", "相声", "脱口秀");
    }
}

监控与告警

/**
 * ES 监控服务
 */
@Slf4j
@Service
public class EsMonitorService {

    @Autowired
    private RestHighLevelClient client;

    @Autowired
    private MeterRegistry meterRegistry;

    /**
     * 1. 记录查询耗时指标
     */
    public <T> List<T> searchWithMetrics(EsQueryBuilder queryBuilder, Class<T> clazz) {
        Timer.Sample sample = Timer.start(meterRegistry);

        try {
            // 执行查询
            List<T> result = esService.search(queryBuilder, clazz);

            // 记录成功指标
            sample.stop(Timer.builder("es.query.duration")
                    .tag("index", queryBuilder.getIndexName())
                    .tag("status", "success")
                    .register(meterRegistry));

            // 记录结果数量
            meterRegistry.gauge("es.query.result.size",
                    Tags.of("index", queryBuilder.getIndexName()),
                    result.size());

            return result;

        } catch (Exception e) {
            // 记录失败指标
            sample.stop(Timer.builder("es.query.duration")
                    .tag("index", queryBuilder.getIndexName())
                    .tag("status", "error")
                    .register(meterRegistry));

            meterRegistry.counter("es.query.error",
                    Tags.of("index", queryBuilder.getIndexName())).increment();

            throw e;
        }
    }

    /**
     * 2. 集群健康检查
     */
    @Scheduled(fixedRate = 60000)
    public void checkClusterHealth() {
        try {
            ClusterHealthRequest request = new ClusterHealthRequest();
            ClusterHealthResponse response = client.cluster()
                    .health(request, RequestOptions.DEFAULT);

            String status = response.getStatus().name();
            int numberOfNodes = response.getNumberOfNodes();
            int activeShards = response.getActiveShards();

            // 记录指标
            meterRegistry.gauge("es.cluster.nodes", numberOfNodes);
            meterRegistry.gauge("es.cluster.shards.active", activeShards);

            // 状态告警
            if ("RED".equals(status)) {
                sendAlert("ES 集群状态异常: RED",
                        String.format("节点数: %d, 活跃分片: %d",
                                numberOfNodes, activeShards));
            } else if ("YELLOW".equals(status)) {
                log.warn("ES 集群状态: YELLOW,可能有副本未分配");
            }

            log.info("ES 集群健康检查: status={}, nodes={}, activeShards={}",
                    status, numberOfNodes, activeShards);

        } catch (Exception e) {
            log.error("ES 集群健康检查失败", e);
            sendAlert("ES 集群健康检查失败", e.getMessage());
        }
    }

    /**
     * 3. 索引状态监控
     */
    @Scheduled(fixedRate = 300000)
    public void checkIndexStats() {
        try {
            IndicesStatsRequest request = new IndicesStatsRequest();
            IndicesStatsResponse response = client.indices()
                    .stats(request, RequestOptions.DEFAULT);

            response.getIndices().forEach((indexName, stats) -> {
                long docCount = stats.getPrimaries().getDocs().getCount();
                long storeSizeBytes = stats.getPrimaries().getStore().getSizeInBytes();

                meterRegistry.gauge("es.index.docs.count",
                        Tags.of("index", indexName), docCount);
                meterRegistry.gauge("es.index.store.size.bytes",
                        Tags.of("index", indexName), storeSizeBytes);

                log.debug("索引 [{}] 统计: 文档数={}, 存储大小={}MB",
                        indexName, docCount, storeSizeBytes / 1024 / 1024);
            });

        } catch (Exception e) {
            log.error("索引状态监控失败", e);
        }
    }

    /**
     * 4. 慢查询日志收集
     */
    public void logSlowQuery(String indexName, String dsl, long elapsed) {
        if (elapsed > 3000) { // 超过3秒
            log.warn("[慢查询] 索引: {}, 耗时: {}ms, DSL: {}",
                    indexName, elapsed, dsl);

            // 发送告警
            if (elapsed > 10000) { // 超过10秒
                sendAlert("ES 慢查询告警",
                        String.format("索引: %s, 耗时: %dms", indexName, elapsed));
            }
        }
    }

    private void sendAlert(String title, String content) {
        // 发送告警通知(钉钉、企业微信、邮件等)
        log.error("【告警】{}: {}", title, content);
    }
}

十一、总结

7-Blog/后端与微服务/assets/Elasticsearch_知识总结-3b295165